14-systemStorageManager.js 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449
  1. /**
  2. * 云函数:systemStorageManager
  3. *
  4. * 通用业务实体、审计和迁移状态的云端入口。
  5. * 数据主存储为 Parse 项目命名空间 Class:
  6. * - VideoWorkflowEntity
  7. * - VideoWorkflowAudit
  8. * - VideoWorkflowMigration
  9. *
  10. * 所有数据归属都通过 _session.js 从 Parse sessionToken 解析得到。
  11. */
  12. const { requireSession, requireAdmin, assertRequestedUserMatchesSession } = require('./_session');
  13. const {
  14. VIDEO_WORKFLOW_CLASSES,
  15. createParseClassStore,
  16. parseUserFields,
  17. parseUserWhere,
  18. } = require('./_parseClassStore');
  19. async function handler(request, response) {
  20. try {
  21. const action = pickParam(request, 'action') || 'list';
  22. if (action === 'adminInspect') {
  23. const adminSession = await requireAdmin(request, Psql);
  24. const store = createParseClassStore({ sessionToken: adminSession.sessionToken });
  25. return adminInspect(request, response, adminSession, store);
  26. }
  27. const session = await requireSession(request, Psql);
  28. assertRequestedUserMatchesSession(request, session);
  29. const store = createParseClassStore({ sessionToken: session.sessionToken });
  30. const userId = session.userId;
  31. if (action === 'list') return listEntities(request, response, userId, store);
  32. if (action === 'get') return getEntity(request, response, userId, store);
  33. if (action === 'upsert') return upsertEntity(request, response, userId, store);
  34. if (action === 'patch') return patchEntity(request, response, userId, store);
  35. if (action === 'delete') return deleteEntity(request, response, userId, store);
  36. if (action === 'purge') return purgeEntity(request, response, userId, store);
  37. if (action === 'audit') return writeAudit(request, response, userId, store);
  38. if (action === 'auditList') return listAudits(request, response, userId, store);
  39. if (action === 'stats') return getStats(request, response, userId, store);
  40. if (action === 'migrationGet') return getMigration(request, response, userId, store);
  41. if (action === 'migrationSet') return setMigration(request, response, userId, store);
  42. return response.json({ code: 400, success: false, error: `未知 action: ${action}` });
  43. } catch (error) {
  44. console.error('systemStorageManager failed:', error && error.message);
  45. return response.json({ code: error.status || 500, success: false, error: error.message });
  46. }
  47. }
  48. async function listEntities(request, response, userId, store) {
  49. const entityType = requiredText(request, 'entityType', 'type');
  50. const status = pickParam(request, 'status') || 'active';
  51. const limit = safeLimit(pickParam(request, 'limit'), 100, 1000);
  52. const where = parseUserWhere(userId, {
  53. entityType,
  54. ...(status ? { status } : {}),
  55. });
  56. const rows = await store.list(VIDEO_WORKFLOW_CLASSES.entity, where, { order: '-updatedAt', limit });
  57. return response.json({ code: 200, success: true, data: rows.map(rowToEntity) });
  58. }
  59. async function getEntity(request, response, userId, store) {
  60. const entityType = requiredText(request, 'entityType', 'type');
  61. const entityId = requiredText(request, 'entityId', 'id', 'bizId');
  62. const row = await findEntity(store, userId, entityType, entityId);
  63. return response.json({ code: 200, success: true, data: row ? rowToEntity(row) : null });
  64. }
  65. async function upsertEntity(request, response, userId, store) {
  66. const entityType = requiredText(request, 'entityType', 'type');
  67. const entityId = requiredText(request, 'entityId', 'id', 'bizId');
  68. const data = normalizeObject(pickParam(request, 'data', 'entity') || {});
  69. // 行状态只表达存储生命周期,业务状态继续放在 data.status,避免 completed/failed 等业务态污染刷新查询。
  70. const status = pickParam(request, 'status') || 'active';
  71. const schemaVersion = Number(pickParam(request, 'schemaVersion') || data.schemaVersion || 1);
  72. const spaceId = pickParam(request, 'spaceId') || data.spaceId || '';
  73. const row = await saveEntity(store, { userId, entityType, entityId, data, status, schemaVersion, spaceId });
  74. return response.json({ code: 200, success: true, data: rowToEntity(row) });
  75. }
  76. async function patchEntity(request, response, userId, store) {
  77. const entityType = requiredText(request, 'entityType', 'type');
  78. const entityId = requiredText(request, 'entityId', 'id', 'bizId');
  79. const patch = normalizeObject(pickParam(request, 'patch', 'data') || {});
  80. const previousRow = await findEntity(store, userId, entityType, entityId);
  81. const previous = previousRow ? normalizeObject(previousRow.data) : {};
  82. const status = pickParam(request, 'status') || previousRow?.status || 'active';
  83. const row = await saveEntity(store, {
  84. userId,
  85. entityType,
  86. entityId,
  87. data: { ...previous, ...patch, id: entityId },
  88. status,
  89. schemaVersion: Number(patch.schemaVersion || previousRow?.schemaVersion || 1),
  90. spaceId: patch.spaceId || previousRow?.spaceId || '',
  91. });
  92. return response.json({ code: 200, success: true, data: rowToEntity(row) });
  93. }
  94. async function deleteEntity(request, response, userId, store) {
  95. const entityType = requiredText(request, 'entityType', 'type');
  96. const entityId = requiredText(request, 'entityId', 'id', 'bizId');
  97. const row = await findEntity(store, userId, entityType, entityId);
  98. if (row?.objectId) {
  99. await store.update(VIDEO_WORKFLOW_CLASSES.entity, row.objectId, { status: 'deleted' });
  100. }
  101. return response.json({ code: 200, success: true, data: { entityType, entityId, deleted: true } });
  102. }
  103. async function purgeEntity(request, response, userId, store) {
  104. const entityType = requiredText(request, 'entityType', 'type');
  105. const entityId = requiredText(request, 'entityId', 'id', 'bizId');
  106. const reason = String(pickParam(request, 'reason') || '').trim();
  107. if (!reason) {
  108. return response.json({ code: 400, success: false, error: '清理必须提供 reason' });
  109. }
  110. const row = await findEntity(store, userId, entityType, entityId);
  111. const purgedAt = new Date().toISOString();
  112. const data = { purged: true, purgedAt, reason, id: entityId, userId };
  113. let purged = false;
  114. if (row?.objectId) {
  115. await store.update(VIDEO_WORKFLOW_CLASSES.entity, row.objectId, {
  116. status: 'purged',
  117. data,
  118. });
  119. purged = true;
  120. }
  121. await createAudit(store, userId, {
  122. entityType,
  123. entityId,
  124. action: 'purge',
  125. summary: `清理 ${entityType}:${entityId}`,
  126. detail: { reason, purged },
  127. });
  128. return response.json({
  129. code: 200,
  130. success: true,
  131. data: { entityType, entityId, purged },
  132. });
  133. }
  134. async function writeAudit(request, response, userId, store) {
  135. const action = requiredText(request, 'auditAction', 'eventAction', 'name');
  136. const entityType = pickParam(request, 'entityType') || '';
  137. const entityId = pickParam(request, 'entityId', 'id') || '';
  138. const summary = String(pickParam(request, 'summary') || '').slice(0, 500);
  139. const detail = normalizeObject(pickParam(request, 'detail') || {});
  140. const row = await createAudit(store, userId, { entityType, entityId, action, summary, detail });
  141. return response.json({ code: 200, success: true, data: rowToAudit(row) });
  142. }
  143. async function listAudits(request, response, userId, store) {
  144. const limit = safeLimit(pickParam(request, 'limit'), 50, 200);
  145. const entityType = String(pickParam(request, 'entityType') || '').trim();
  146. const where = parseUserWhere(userId, entityType ? { entityType } : {});
  147. const rows = await store.list(VIDEO_WORKFLOW_CLASSES.audit, where, { order: '-createdAt', limit });
  148. return response.json({ code: 200, success: true, data: rows.map(rowToAudit) });
  149. }
  150. async function getStats(request, response, userId, store) {
  151. const [entityRows, auditRows, migrationRows] = await Promise.all([
  152. listAll(store, VIDEO_WORKFLOW_CLASSES.entity, parseUserWhere(userId, {}), 5000),
  153. listAll(store, VIDEO_WORKFLOW_CLASSES.audit, parseUserWhere(userId, {}), 1000),
  154. listAll(store, VIDEO_WORKFLOW_CLASSES.migration, parseUserWhere(userId, {}), 1000),
  155. ]);
  156. return response.json({
  157. code: 200,
  158. success: true,
  159. data: buildStats(entityRows, auditRows, migrationRows),
  160. });
  161. }
  162. async function adminInspect(request, response, adminSession, store) {
  163. const targetUserId = String(pickParam(request, 'targetUserId', 'inspectUserId') || '').trim();
  164. const limit = safeLimit(pickParam(request, 'limit'), 100, 500);
  165. const userWhere = targetUserId ? { ownerId: targetUserId } : {};
  166. const [entityRows, auditRows, migrationRows] = await Promise.all([
  167. listAll(store, VIDEO_WORKFLOW_CLASSES.entity, parseUserWhereForAdmin(userWhere), limit * 20),
  168. listAll(store, VIDEO_WORKFLOW_CLASSES.audit, parseUserWhereForAdmin(userWhere), limit * 20),
  169. listAll(store, VIDEO_WORKFLOW_CLASSES.migration, parseUserWhereForAdmin(userWhere), limit * 20),
  170. ]);
  171. const users = {};
  172. const ensureUser = (userId) => {
  173. if (!users[userId]) {
  174. users[userId] = {
  175. userId,
  176. totalEntities: 0,
  177. entityTypes: {},
  178. auditCount: 0,
  179. latestAuditAt: '',
  180. migrations: { total: 0, failedRows: 0, byStatus: {} },
  181. latestUpdatedAt: '',
  182. };
  183. }
  184. return users[userId];
  185. };
  186. for (const row of entityRows) {
  187. const user = ensureUser(row.ownerId || '');
  188. const entityType = row.entityType || '';
  189. const status = row.status || 'active';
  190. if (!entityType || !user.userId) continue;
  191. if (!user.entityTypes[entityType]) user.entityTypes[entityType] = { total: 0, byStatus: {} };
  192. user.entityTypes[entityType].total += 1;
  193. user.entityTypes[entityType].byStatus[status] = (user.entityTypes[entityType].byStatus[status] || 0) + 1;
  194. user.totalEntities += 1;
  195. const latest = row.updatedAt || '';
  196. if (latest && String(latest) > String(user.latestUpdatedAt || '')) user.latestUpdatedAt = latest;
  197. }
  198. for (const row of auditRows) {
  199. const user = ensureUser(row.ownerId || '');
  200. if (!user.userId) continue;
  201. user.auditCount += 1;
  202. const latest = row.createdAt || '';
  203. if (latest && String(latest) > String(user.latestAuditAt || '')) user.latestAuditAt = latest;
  204. }
  205. for (const row of migrationRows) {
  206. const user = ensureUser(row.ownerId || '');
  207. if (!user.userId) continue;
  208. const status = row.status || 'unknown';
  209. user.migrations.total += 1;
  210. user.migrations.byStatus[status] = (user.migrations.byStatus[status] || 0) + 1;
  211. user.migrations.failedRows += Number(row.failedCount || 0);
  212. }
  213. const rows = Object.values(users)
  214. .sort((a, b) => b.totalEntities - a.totalEntities || String(b.latestUpdatedAt || '').localeCompare(String(a.latestUpdatedAt || '')))
  215. .slice(0, limit);
  216. return response.json({
  217. code: 200,
  218. success: true,
  219. data: {
  220. adminUserId: adminSession.userId,
  221. targetUserId,
  222. users: rows,
  223. generatedAt: new Date().toISOString(),
  224. },
  225. });
  226. }
  227. async function getMigration(request, response, userId, store) {
  228. const migrationKey = requiredText(request, 'migrationKey');
  229. const sourceKey = requiredText(request, 'sourceKey');
  230. const row = await findMigration(store, userId, migrationKey, sourceKey);
  231. return response.json({ code: 200, success: true, data: row ? rowToMigration(row) : null });
  232. }
  233. async function setMigration(request, response, userId, store) {
  234. const migrationKey = requiredText(request, 'migrationKey');
  235. const sourceKey = requiredText(request, 'sourceKey');
  236. const status = pickParam(request, 'status') || 'completed';
  237. const migratedCount = Number(pickParam(request, 'migratedCount') || 0);
  238. const failedCount = Number(pickParam(request, 'failedCount') || 0);
  239. const detail = normalizeObject(pickParam(request, 'detail') || {});
  240. const where = parseUserWhere(userId, { migrationKey, sourceKey });
  241. const payload = {
  242. ...parseUserFields(userId),
  243. migrationKey,
  244. sourceKey,
  245. status,
  246. migratedCount,
  247. failedCount,
  248. detail,
  249. };
  250. const row = await upsertMigrationWithDetailFallback(store, where, payload);
  251. return response.json({ code: 200, success: true, data: rowToMigration(row) });
  252. }
  253. async function upsertMigrationWithDetailFallback(store, where, payload) {
  254. try {
  255. return await store.upsertByQuery(VIDEO_WORKFLOW_CLASSES.migration, where, payload);
  256. } catch (error) {
  257. if (!isParseFieldTypeMismatch(error) || typeof payload.detail === 'string') {
  258. throw error;
  259. }
  260. return store.upsertByQuery(VIDEO_WORKFLOW_CLASSES.migration, where, {
  261. ...payload,
  262. detail: safeStringify(payload.detail),
  263. });
  264. }
  265. }
  266. async function saveEntity(store, { userId, entityType, entityId, data, status, schemaVersion, spaceId }) {
  267. const rawData = normalizeObject(data);
  268. const payload = { ...rawData, id: rawData.id || entityId, userId };
  269. const where = parseUserWhere(userId, { entityType, entityId });
  270. return store.upsertByQuery(VIDEO_WORKFLOW_CLASSES.entity, where, {
  271. ...parseUserFields(userId),
  272. entityType,
  273. entityId,
  274. status: status || 'active',
  275. schemaVersion: schemaVersion || 1,
  276. spaceId: spaceId || '',
  277. data: payload,
  278. });
  279. }
  280. async function findEntity(store, userId, entityType, entityId) {
  281. return store.findFirst(VIDEO_WORKFLOW_CLASSES.entity, parseUserWhere(userId, { entityType, entityId }), { limit: 1 });
  282. }
  283. async function findMigration(store, userId, migrationKey, sourceKey) {
  284. return store.findFirst(VIDEO_WORKFLOW_CLASSES.migration, parseUserWhere(userId, { migrationKey, sourceKey }), { limit: 1 });
  285. }
  286. async function createAudit(store, userId, { entityType, entityId, action, summary, detail }) {
  287. return store.create(VIDEO_WORKFLOW_CLASSES.audit, {
  288. ...parseUserFields(userId),
  289. entityType: entityType || '',
  290. entityId: entityId || '',
  291. action,
  292. summary: String(summary || '').slice(0, 500),
  293. detail: safeStringify(detail || {}),
  294. });
  295. }
  296. async function listAll(store, className, where, maxRows) {
  297. const all = [];
  298. const pageSize = 1000;
  299. for (let skip = 0; skip < maxRows; skip += pageSize) {
  300. const rows = await store.list(className, where, { order: '-updatedAt', limit: Math.min(pageSize, maxRows - skip), skip });
  301. all.push(...rows);
  302. if (rows.length < pageSize || all.length >= maxRows) break;
  303. }
  304. return all.slice(0, maxRows);
  305. }
  306. function buildStats(entityRows, auditRows, migrationRows) {
  307. const entityTypes = {};
  308. let totalEntities = 0;
  309. for (const row of entityRows) {
  310. const entityType = row.entityType || '';
  311. const status = row.status || 'active';
  312. if (!entityType) continue;
  313. if (!entityTypes[entityType]) entityTypes[entityType] = { total: 0, byStatus: {} };
  314. entityTypes[entityType].total += 1;
  315. entityTypes[entityType].byStatus[status] = (entityTypes[entityType].byStatus[status] || 0) + 1;
  316. totalEntities += 1;
  317. }
  318. const migrations = { total: 0, byStatus: {} };
  319. for (const row of migrationRows) {
  320. const status = row.status || 'unknown';
  321. migrations.total += 1;
  322. migrations.byStatus[status] = (migrations.byStatus[status] || 0) + 1;
  323. }
  324. return {
  325. totalEntities,
  326. entityTypes,
  327. auditCount: auditRows.length,
  328. migrations,
  329. generatedAt: new Date().toISOString(),
  330. };
  331. }
  332. function parseUserWhereForAdmin(extra = {}) {
  333. return { projectKey: 'video-workflow', ...extra };
  334. }
  335. function rowToEntity(row) {
  336. const data = normalizeObject(row.data);
  337. return {
  338. id: row.entityId,
  339. entityId: row.entityId,
  340. type: row.entityType,
  341. entityType: row.entityType,
  342. status: row.status || 'active',
  343. schemaVersion: row.schemaVersion,
  344. spaceId: row.spaceId || '',
  345. data,
  346. createdAt: row.createdAt,
  347. updatedAt: row.updatedAt,
  348. };
  349. }
  350. function rowToMigration(row) {
  351. return {
  352. id: row.objectId,
  353. migrationKey: row.migrationKey,
  354. sourceKey: row.sourceKey,
  355. status: row.status,
  356. migratedCount: Number(row.migratedCount || 0),
  357. failedCount: Number(row.failedCount || 0),
  358. detail: normalizeObject(row.detail),
  359. createdAt: row.createdAt,
  360. updatedAt: row.updatedAt,
  361. };
  362. }
  363. function rowToAudit(row) {
  364. return {
  365. id: row.objectId,
  366. entityType: row.entityType || '',
  367. entityId: row.entityId || '',
  368. action: row.action,
  369. summary: row.summary || '',
  370. detail: normalizeObject(row.detail),
  371. createdAt: row.createdAt,
  372. };
  373. }
  374. function requiredText(request, ...names) {
  375. const value = pickParam(request, ...names);
  376. const text = String(value || '').trim();
  377. if (!text) {
  378. const error = new Error(`缺少 ${names[0]}`);
  379. error.status = 400;
  380. throw error;
  381. }
  382. return text;
  383. }
  384. function pickParam(request, ...names) {
  385. const sources = [request.params, request.body, request.query, request];
  386. for (const src of sources) {
  387. if (!src || typeof src !== 'object') continue;
  388. for (const name of names) {
  389. const value = src[name];
  390. if (value !== undefined && value !== null && value !== '') return value;
  391. }
  392. }
  393. return null;
  394. }
  395. function normalizeObject(value) {
  396. if (typeof value === 'string') {
  397. try { return JSON.parse(value); } catch { return {}; }
  398. }
  399. if (!value || typeof value !== 'object' || Array.isArray(value)) return {};
  400. return value;
  401. }
  402. function safeLimit(value, fallback, max) {
  403. const parsed = parseInt(value || fallback, 10);
  404. return Number.isFinite(parsed) ? Math.min(Math.max(parsed, 1), max) : fallback;
  405. }
  406. function safeStringify(value) {
  407. const text = JSON.stringify(normalizeObject(value));
  408. return text.length > 8000 ? `${text.slice(0, 7997)}...` : text;
  409. }
  410. function isParseFieldTypeMismatch(error) {
  411. const message = String(error?.message || error?.detail?.error || '').toLowerCase();
  412. return Number(error?.status || 0) === 400
  413. && (message.includes('schema') || message.includes('type') || message.includes('expected'));
  414. }