_parseClassStore.js 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248
  1. /**
  2. * 云函数共享的 Parse Class 存储工具。
  3. *
  4. * 本项目的结构化业务主数据写入带 VideoWorkflow 命名空间的 Parse Class。
  5. * 这里封装 REST 调用、项目命名空间字段和 owner Pointer,避免业务云函数
  6. * 继续直接依赖 PSQL。
  7. */
  8. const PARSE_CLASS_PROJECT_KEY = 'video-workflow';
  9. const VIDEO_WORKFLOW_CLASSES = {
  10. entity: 'VideoWorkflowEntity',
  11. audit: 'VideoWorkflowAudit',
  12. migration: 'VideoWorkflowMigration',
  13. fileAsset: 'VideoWorkflowFileAsset',
  14. };
  15. function createParseClassStore(options = {}) {
  16. const parseApiHost = normalizeParseApiHost(options.parseApiHost || parseStoreEnv('PARSE_API_HOST', parseStoreEnv('PARSE_BASE_URL', 'https://server.fmode.cn')));
  17. const appId = options.appId || parseStoreEnv('PARSE_APP_ID', parseStoreEnv('PARSE_APPLICATION_ID', 'ncloudmaster'));
  18. const sessionToken = options.sessionToken || '';
  19. return {
  20. findFirst: (className, where, opts = {}) => parseFindFirst({ parseApiHost, appId, sessionToken, className, where, opts }),
  21. list: (className, where, opts = {}) => parseList({ parseApiHost, appId, sessionToken, className, where, opts }),
  22. create: (className, data) => parseCreate({ parseApiHost, appId, sessionToken, className, data }),
  23. update: (className, objectId, patch) => parseUpdate({ parseApiHost, appId, sessionToken, className, objectId, patch }),
  24. deleteRecord: (className, objectId) => parseDelete({ parseApiHost, appId, sessionToken, className, objectId }),
  25. upsertByQuery: (className, where, data, opts = {}) => parseUpsertByQuery({ parseApiHost, appId, sessionToken, className, where, data, opts }),
  26. userFields: (userId) => parseUserFields(userId),
  27. projectKey: PARSE_CLASS_PROJECT_KEY,
  28. };
  29. }
  30. function parseUserFields(userId) {
  31. return {
  32. projectKey: PARSE_CLASS_PROJECT_KEY,
  33. owner: { __type: 'Pointer', className: '_User', objectId: String(userId || '') },
  34. ownerId: String(userId || ''),
  35. };
  36. }
  37. function parseUserWhere(userId, extra = {}) {
  38. return {
  39. projectKey: PARSE_CLASS_PROJECT_KEY,
  40. owner: { __type: 'Pointer', className: '_User', objectId: String(userId || '') },
  41. ...extra,
  42. };
  43. }
  44. async function parseFindFirst(ctx) {
  45. const rows = await parseList({ ...ctx, opts: { ...ctx.opts, limit: 1 } });
  46. return rows[0] || null;
  47. }
  48. async function parseList({ parseApiHost, appId, sessionToken, className, where, opts = {} }) {
  49. const query = new URLSearchParams();
  50. query.set('where', JSON.stringify(where || {}));
  51. query.set('limit', String(safeParseLimit(opts.limit, 100, 1000)));
  52. if (opts.skip !== undefined) query.set('skip', String(Math.max(0, parseInt(opts.skip || 0, 10) || 0)));
  53. if (opts.order) query.set('order', String(opts.order));
  54. if (opts.include) query.set('include', String(opts.include));
  55. if (opts.keys) query.set('keys', String(opts.keys));
  56. const data = await parseRequest({
  57. parseApiHost,
  58. appId,
  59. sessionToken,
  60. method: 'GET',
  61. path: `/classes/${encodeURIComponent(className)}?${query.toString()}`,
  62. });
  63. return Array.isArray(data.results) ? data.results : [];
  64. }
  65. async function parseCreate({ parseApiHost, appId, sessionToken, className, data }) {
  66. const created = await parseRequest({
  67. parseApiHost,
  68. appId,
  69. sessionToken,
  70. method: 'POST',
  71. path: `/classes/${encodeURIComponent(className)}`,
  72. body: sanitizeParseWrite(data),
  73. });
  74. return { ...sanitizeParseWrite(data), ...created };
  75. }
  76. async function parseUpdate({ parseApiHost, appId, sessionToken, className, objectId, patch }) {
  77. await parseRequest({
  78. parseApiHost,
  79. appId,
  80. sessionToken,
  81. method: 'PUT',
  82. path: `/classes/${encodeURIComponent(className)}/${encodeURIComponent(objectId)}`,
  83. body: sanitizeParseWrite(patch),
  84. });
  85. return parseFindFirst({
  86. parseApiHost,
  87. appId,
  88. sessionToken,
  89. className,
  90. where: { objectId },
  91. opts: { limit: 1 },
  92. });
  93. }
  94. async function parseDelete({ parseApiHost, appId, sessionToken, className, objectId }) {
  95. return parseRequest({
  96. parseApiHost,
  97. appId,
  98. sessionToken,
  99. method: 'DELETE',
  100. path: `/classes/${encodeURIComponent(className)}/${encodeURIComponent(objectId)}`,
  101. });
  102. }
  103. async function parseUpsertByQuery({ parseApiHost, appId, sessionToken, className, where, data, opts = {} }) {
  104. const findExisting = () => parseFindFirst({
  105. parseApiHost,
  106. appId,
  107. sessionToken,
  108. className,
  109. where,
  110. opts: { order: opts.order || '-updatedAt', limit: 1 },
  111. });
  112. const existing = await findExisting();
  113. if (existing && existing.objectId) {
  114. return parseUpdate({
  115. parseApiHost,
  116. appId,
  117. sessionToken,
  118. className,
  119. objectId: existing.objectId,
  120. patch: data,
  121. });
  122. }
  123. try {
  124. return await parseCreate({ parseApiHost, appId, sessionToken, className, data: { ...where, ...data } });
  125. } catch (error) {
  126. if (!isRetryableParseError(error)) throw error;
  127. await parseStoreSleep(500);
  128. const created = await findExisting();
  129. if (created && created.objectId) {
  130. return parseUpdate({
  131. parseApiHost,
  132. appId,
  133. sessionToken,
  134. className,
  135. objectId: created.objectId,
  136. patch: data,
  137. });
  138. }
  139. return parseCreate({ parseApiHost, appId, sessionToken, className, data: { ...where, ...data } });
  140. }
  141. }
  142. async function parseRequest({ parseApiHost, appId, sessionToken, method, path, body }) {
  143. const headers = {
  144. 'Accept': 'application/json',
  145. 'X-Parse-Application-Id': appId,
  146. };
  147. if (sessionToken) headers['X-Parse-Session-Token'] = sessionToken;
  148. if (body !== undefined) headers['Content-Type'] = 'application/json';
  149. const maxAttempts = isSafeParseRetryMethod(method) ? 3 : 1;
  150. let lastError = null;
  151. for (let attempt = 1; attempt <= maxAttempts; attempt += 1) {
  152. try {
  153. const resp = await fetch(`${parseApiHost}/parse${path}`, {
  154. method,
  155. headers,
  156. body: body !== undefined ? JSON.stringify(body) : undefined,
  157. });
  158. const data = await resp.json().catch(() => ({}));
  159. if (!resp.ok || data.error) {
  160. const error = new Error(data.error || data.message || `Parse ${method} ${path} failed`);
  161. error.status = resp.status || 500;
  162. error.detail = data;
  163. if (attempt < maxAttempts && isRetryableParseError(error)) {
  164. await parseStoreSleep(350 * attempt);
  165. continue;
  166. }
  167. throw error;
  168. }
  169. return data;
  170. } catch (error) {
  171. lastError = error;
  172. if (attempt < maxAttempts && isRetryableParseError(error)) {
  173. await parseStoreSleep(350 * attempt);
  174. continue;
  175. }
  176. throw error;
  177. }
  178. }
  179. throw lastError || new Error(`Parse ${method} ${path} failed`);
  180. }
  181. function sanitizeParseWrite(value) {
  182. const source = normalizeParseObject(value);
  183. const next = {};
  184. for (const [key, item] of Object.entries(source)) {
  185. if (['objectId', 'createdAt', 'updatedAt'].includes(key)) continue;
  186. if (item === undefined) continue;
  187. next[key] = item;
  188. }
  189. return next;
  190. }
  191. function normalizeParseObject(value) {
  192. if (!value || typeof value !== 'object' || Array.isArray(value)) return {};
  193. return value;
  194. }
  195. function normalizeParseApiHost(value) {
  196. return String(value || 'https://server.fmode.cn').replace(/\/+$/, '').replace(/\/parse$/i, '');
  197. }
  198. function safeParseLimit(value, fallback, max) {
  199. const parsed = parseInt(value || fallback, 10);
  200. return Number.isFinite(parsed) ? Math.min(Math.max(parsed, 1), max) : fallback;
  201. }
  202. function parseStoreEnv(name, fallback = '') {
  203. if (typeof process !== 'undefined' && process.env && process.env[name] !== undefined) {
  204. return process.env[name];
  205. }
  206. return fallback;
  207. }
  208. function isSafeParseRetryMethod(method) {
  209. return ['GET', 'PUT', 'DELETE'].includes(String(method || '').toUpperCase());
  210. }
  211. function isRetryableParseError(error) {
  212. const status = Number(error?.status || 0);
  213. const message = `${error?.message || ''} ${error?.detail?.error || ''} ${error?.detail?.message || ''}`;
  214. return status >= 500 || /fetch failed|Failed to fetch|NetworkError|Load failed|ECONNRESET|ETIMEDOUT|EAI_AGAIN/i.test(message);
  215. }
  216. function parseStoreSleep(ms) {
  217. return new Promise((resolve) => setTimeout(resolve, ms));
  218. }
  219. module.exports = {
  220. PARSE_CLASS_PROJECT_KEY,
  221. VIDEO_WORKFLOW_CLASSES,
  222. createParseClassStore,
  223. parseUserFields,
  224. parseUserWhere,
  225. };