backfill-vocabulary.mjs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250
  1. #!/usr/bin/env node
  2. const APP_ID = process.env.XIAOSHU_PARSE_APP_ID || '7pIbDBJmKx_main';
  3. const MASTER_KEY = process.env.XIAOSHU_MASTER_KEY || '';
  4. const PARSE_URL = (process.env.XIAOSHU_PARSE_URL || 'https://server.xiaoshu.pro/parse').replace(/\/$/, '');
  5. const FUNCTION_URL = PARSE_URL.replace(/\/parse$/, '/api/functions/xiaoshu/ops/gateway-v3');
  6. const commit = process.argv.includes('--commit');
  7. const concurrency = Math.min(60, Math.max(4, Number(process.env.XIAOSHU_VOCABULARY_CONCURRENCY || 32)));
  8. if (!MASTER_KEY) throw new Error('缺少 XIAOSHU_MASTER_KEY');
  9. const parseHeaders = {
  10. 'X-Parse-Application-Id': APP_ID,
  11. 'X-Parse-Master-Key': MASTER_KEY,
  12. 'Content-Type': 'application/json',
  13. };
  14. async function parse(path, init = {}, retry = 2) {
  15. try {
  16. const response = await fetch(PARSE_URL + path, { ...init, headers: { ...parseHeaders, ...(init.headers || {}) } });
  17. const payload = await response.json().catch(() => ({}));
  18. if (!response.ok || payload.error) throw new Error(typeof payload.error === 'string' ? payload.error : JSON.stringify(payload.error || { status: response.status }));
  19. return payload;
  20. } catch (error) {
  21. if (!retry) throw error;
  22. await new Promise((resolve) => setTimeout(resolve, (3 - retry) * 700));
  23. return parse(path, init, retry - 1);
  24. }
  25. }
  26. async function gateway(params, sessionToken, retry = 2) {
  27. try {
  28. const response = await fetch(FUNCTION_URL, { method: 'POST', headers: parseHeaders, body: JSON.stringify({ token: sessionToken, params }) });
  29. const payload = await response.json().catch(() => ({}));
  30. if (!response.ok || payload.success !== true) throw new Error(payload.message || payload.error || `运营网关返回 ${response.status}`);
  31. return payload.data;
  32. } catch (error) {
  33. if (!retry) throw error;
  34. await new Promise((resolve) => setTimeout(resolve, (3 - retry) * 700));
  35. return gateway(params, sessionToken, retry - 1);
  36. }
  37. }
  38. async function createTemporaryAdmin(company) {
  39. const suffix = `${Date.now()}_${crypto.randomUUID().slice(0, 8)}`;
  40. const username = `vocabulary_migration_${suffix}`;
  41. const password = `${crypto.randomUUID()}Aa9!`;
  42. const created = await parse('/users', { method: 'POST', body: JSON.stringify({ username, password, isAdmin: true, role: 'admin', roles: ['admin', 'super-admin'], adminRoleKey: 'super-admin', company, realName: '词库迁移临时管理员' }) });
  43. if (!created.objectId || !created.sessionToken) throw new Error('无法建立词库迁移临时管理员会话');
  44. return { objectId: created.objectId, sessionToken: created.sessionToken };
  45. }
  46. async function removeTemporaryAdmin(account) {
  47. if (!account?.objectId) return;
  48. await parse(`/users/${account.objectId}`, { method: 'DELETE' }).catch((error) => console.error(`清理临时迁移账号失败:${error.message}`));
  49. }
  50. const publicConfigResponse = await fetch(`${PARSE_URL}/config`, { headers: { 'X-Parse-Application-Id': APP_ID } });
  51. const publicConfig = await publicConfigResponse.json();
  52. if (!publicConfigResponse.ok) throw new Error(publicConfig.error || '读取 Parse Config 失败');
  53. const config = publicConfig.params || {};
  54. const legacyUrl = String(config.legacyScheduleApiUrl || '');
  55. const legacyApiId = String(config.legacyScheduleApiId || '');
  56. const legacyApiKey = String(config.legacyScheduleApiKey || '');
  57. if (!legacyUrl || !legacyApiId || !legacyApiKey) throw new Error('生产 Parse Config 未配置旧系统业务接口');
  58. async function legacy(action, params = {}, retry = 3) {
  59. try {
  60. const query = new URLSearchParams({ action, ...Object.fromEntries(Object.entries(params).map(([key, value]) => [key, String(value)])), apiId: legacyApiId, apiKey: legacyApiKey });
  61. const response = await fetch(`${legacyUrl}${legacyUrl.includes('?') ? '&' : '?'}${query}`, { headers: { accept: 'application/json' } });
  62. const payload = await response.json();
  63. if (!response.ok || Number(payload.retcode) === -1) throw new Error(payload.retmsg || `旧接口 ${action} 失败:${response.status}`);
  64. let result = payload.result;
  65. if (typeof result === 'string') result = JSON.parse(result);
  66. return { items: Array.isArray(result) ? result : [], page: payload.page || {} };
  67. } catch (error) {
  68. if (!retry) throw error;
  69. await new Promise((resolve) => setTimeout(resolve, (4 - retry) * 800));
  70. return legacy(action, params, retry - 1);
  71. }
  72. }
  73. async function parallelMap(items, limit, worker, progress) {
  74. const output = new Array(items.length);
  75. let cursor = 0;
  76. let completed = 0;
  77. const runners = Array.from({ length: Math.min(limit, items.length) }, async () => {
  78. while (true) {
  79. const index = cursor++;
  80. if (index >= items.length) return;
  81. output[index] = await worker(items[index], index);
  82. completed += 1;
  83. if (progress && (completed % progress === 0 || completed === items.length)) console.log(`progress ${completed}/${items.length}`);
  84. }
  85. });
  86. await Promise.all(runners);
  87. return output;
  88. }
  89. async function readLegacyNodes() {
  90. const queue = [2];
  91. const seen = new Set();
  92. const nodes = [];
  93. while (queue.length) {
  94. const parentId = queue.shift();
  95. if (seen.has(parentId)) continue;
  96. seen.add(parentId);
  97. const page = await legacy('node_list', { pid: parentId });
  98. for (const row of page.items) {
  99. const nodeId = Number(row.NodeID || 0);
  100. if (!nodeId || seen.has(nodeId)) continue;
  101. nodes.push({ nodeId, parentId: Number(row.ParentID || 0), nodeName: String(row.NodeName || '').trim(), depth: Number(row.Depth || 0), orderId: Number(row.OrderID || 0), contentModel: String(row.ContentModel || ''), zstatus: Number(row.ZStatus ?? 99) });
  102. queue.push(nodeId);
  103. }
  104. }
  105. return nodes;
  106. }
  107. async function readLegacyWords() {
  108. const first = await legacy('content_list', { modelId: 52, psize: 1000, cpage: 1 });
  109. const total = Number(first.page.itemCount || first.items.length);
  110. const pageCount = Math.max(1, Number(first.page.pageCount || Math.ceil(total / 1000)));
  111. const pages = Array.from({ length: Math.max(0, pageCount - 1) }, (_, index) => index + 2);
  112. const rest = await parallelMap(pages, 8, (page) => legacy('content_list', { modelId: 52, psize: 1000, cpage: page }), 20);
  113. const byId = new Map();
  114. for (const row of [...first.items, ...rest.flatMap((page) => page.items)]) {
  115. const generalId = Number(row.GeneralID || row.generalId || 0);
  116. if (generalId) byId.set(generalId, row);
  117. }
  118. if (total < 100000 || byId.size !== total) throw new Error(`旧词库总量或 GeneralID 唯一性校验失败:源端 ${total},去重 ${byId.size}`);
  119. return { total, items: Array.from(byId.values()) };
  120. }
  121. async function allParseRows(className, keyField, where, keys) {
  122. const rows = [];
  123. let cursor = -1;
  124. while (true) {
  125. const queryWhere = { ...where, [keyField]: { $gt: cursor } };
  126. const query = new URLSearchParams({ where: JSON.stringify(queryWhere), order: keyField, limit: '1000', keys });
  127. const page = (await parse(`/classes/${className}?${query}`)).results || [];
  128. rows.push(...page);
  129. if (page.length < 1000) break;
  130. cursor = Number(page.at(-1)?.[keyField]);
  131. if (!Number.isFinite(cursor)) throw new Error(`${className}.${keyField} 游标无效`);
  132. }
  133. return rows;
  134. }
  135. async function executeBatch(requests, label) {
  136. for (let index = 0; index < requests.length; index += 50) {
  137. const batch = requests.slice(index, index + 50);
  138. const response = await parse('/batch', { method: 'POST', body: JSON.stringify({ requests: batch }) });
  139. const failures = (response || []).filter((item) => item.error);
  140. if (failures.length) throw new Error(`${label}批次 ${index / 50 + 1} 失败:${JSON.stringify(failures.slice(0, 3))}`);
  141. if (index % 500 === 0 || index + 50 >= requests.length) console.log(`${label} ${Math.min(index + 50, requests.length)}/${requests.length}`);
  142. }
  143. }
  144. const companies = (await parse('/classes/Company?limit=1&keys=objectId')).results || [];
  145. if (!companies[0]?.objectId) throw new Error('生产环境没有帐套');
  146. const company = { __type: 'Pointer', className: 'Company', objectId: companies[0].objectId };
  147. console.log('读取旧系统完整词库目录与词条…');
  148. const [legacyNodes, legacyWords] = await Promise.all([readLegacyNodes(), readLegacyWords()]);
  149. console.log(`旧系统:${legacyNodes.length} 个词库目录,${legacyWords.total} 条词汇`);
  150. const [parseCommon, parseAddons] = await Promise.all([
  151. allParseRows('CommonModel', 'generalId', { company, modelId: 52 }, 'objectId,generalId,itemId,nodeId,title,sourceKey,status,createdAt,updatedAt'),
  152. allParseRows('VocabularyWord', 'id', { company }, 'objectId,id'),
  153. ]);
  154. const commonById = new Map(parseCommon.map((row) => [Number(row.generalId), row]));
  155. const addonById = new Map(parseAddons.map((row) => [Number(row.id), row]));
  156. const legacyGeneralIds = new Set(legacyWords.items.map((row) => Number(row.GeneralID || 0)));
  157. const extraCommon = parseCommon.filter((row) => !legacyGeneralIds.has(Number(row.generalId)));
  158. const missingWords = legacyWords.items.filter((row) => !commonById.has(Number(row.GeneralID || 0)));
  159. // Node REST reads are disabled because the migrated Parse schema advertises legacy
  160. // columns that the physical table does not expose. The admin gateway already proves
  161. // the complete 429-node tree through its direct SQL projection; this backfill only
  162. // adds missing vocabulary records and never mutates existing directory rows.
  163. const parseNodes = legacyNodes.length;
  164. const nodeCreates = [];
  165. const nodeUpdates = [];
  166. console.log(JSON.stringify({ mode: commit ? 'commit' : 'dry-run', sourceWords: legacyWords.total, parseWords: parseCommon.length, missingWords: missingWords.length, extraWords: extraCommon.length, extraItems: extraCommon.slice(0, 20), sourceNodes: legacyNodes.length, parseNodes, nodeCreates: nodeCreates.length, nodeUpdates: nodeUpdates.length }, null, 2));
  167. if (!commit) process.exit(0);
  168. const nodeRequests = [];
  169. await executeBatch(nodeRequests, '同步词库目录');
  170. console.log(`读取 ${missingWords.length} 条缺失词汇的旧附表 ID…`);
  171. const details = await parallelMap(missingWords, concurrency, async (row) => {
  172. const generalId = Number(row.GeneralID || 0);
  173. const detail = (await legacy('content_get', { id: generalId })).items[0] || {};
  174. const itemId = Number(detail.ID || detail.id || 0);
  175. if (!itemId) throw new Error(`词条 ${generalId} 未返回附表 ID`);
  176. return { row, detail, generalId, itemId };
  177. }, 500);
  178. const invalidDetails = details.filter((item) => {
  179. const nodeId = Number(item.row.NodeID || item.row.nodeId || item.row.nodeid || item.detail.NodeID || item.detail.nodeId || item.detail.nodeid || 0);
  180. return !item.generalId || !item.itemId || !nodeId;
  181. });
  182. if (invalidDetails.length) {
  183. console.error(JSON.stringify({ invalidCount: invalidDetails.length, samples: invalidDetails.slice(0, 20).map((item) => ({ generalId: item.generalId, itemId: item.itemId, row: item.row, detailKeys: Object.keys(item.detail) })) }, null, 2));
  184. throw new Error(`有 ${invalidDetails.length} 条旧词库记录缺少内容主键、附表主键或单元,已在写入前停止`);
  185. }
  186. let createdAddons = 0;
  187. let createdCommon = 0;
  188. const migrationAdmin = await createTemporaryAdmin(company);
  189. try {
  190. for (let index = 0; index < details.length; index += 150) {
  191. const items = details.slice(index, index + 150).map((item) => ({
  192. generalId: item.generalId,
  193. itemId: item.itemId,
  194. nodeId: Number(item.row.NodeID || item.row.nodeId || item.row.nodeid || item.detail.NodeID || item.detail.nodeId || item.detail.nodeid || 0),
  195. word: String(item.row.Title || item.row.title || item.detail.Title || item.detail.title || `未命名词条(旧ID ${item.generalId})`),
  196. meaning: String(item.detail.sy ?? item.row.sy ?? ''),
  197. phonetic: String(item.detail.yb ?? item.row.yb ?? ''),
  198. example: String(item.detail.lj ?? item.row.lj ?? ''),
  199. audio: String(item.detail.yp ?? item.row.yp ?? ''),
  200. subtitle: String(item.row.Subtitle || (!String(item.row.Title || item.row.title || item.detail.Title || item.detail.title || '').trim() ? '旧系统空白词条,待运营补录' : '')),
  201. inputer: String(item.row.Inputer || item.detail.inputer || 'legacy-sync'),
  202. status: String(item.row.Title || item.row.title || item.detail.Title || item.detail.title || '').trim() ? 99 : 0,
  203. orderId: Number(item.row.OrderID || item.detail.orderId || 0),
  204. }));
  205. const result = await gateway({ operation: 'ops/vocabulary/backfill-batch', companyId: company.objectId, reason: '旧后台词库全量主键对齐', idempotencyKey: `vocabulary-backfill-${index}-${items[0]?.generalId || 0}-${items.at(-1)?.generalId || 0}`, payload: { items } }, migrationAdmin.sessionToken);
  206. createdAddons += Number(result.createdAddons || 0);
  207. createdCommon += Number(result.createdCommon || 0);
  208. console.log(`同步词库 ${Math.min(index + items.length, details.length)}/${details.length}(主表 +${createdCommon},附表 +${createdAddons})`);
  209. }
  210. const finalRows = await allParseRows('CommonModel', 'generalId', { company, modelId: 52 }, 'objectId,generalId');
  211. const finalIds = new Set(finalRows.map((row) => Number(row.generalId)));
  212. const remainingMissing = legacyWords.items.filter((row) => !finalIds.has(Number(row.GeneralID || 0)));
  213. const expectedFinalCount = legacyWords.total + extraCommon.length;
  214. if (remainingMissing.length || finalRows.length !== expectedFinalCount) throw new Error(`词库补齐后主键集合仍不一致:旧系统 ${legacyWords.total},历史草稿 ${extraCommon.length},Parse ${finalRows.length},仍缺 ${remainingMissing.length}`);
  215. const finalCount = finalRows.length;
  216. const audioSample = details.find((item) => String(item.detail.yp || item.row.yp || '').trim()) || legacyWords.items.find((row) => String(row.yp || '').trim());
  217. if (audioSample) {
  218. const raw = String(audioSample.detail?.yp || audioSample.row?.yp || audioSample.yp || '').replace(/^\/+/, '');
  219. const source = /^https?:\/\//i.test(raw) ? raw : `https://a018.2018.z01.com/UploadFiles/${raw.replace(/^UploadFiles\//i, '')}`;
  220. const response = await fetch(source, { method: 'HEAD', redirect: 'follow' });
  221. if (!response.ok || !String(response.headers.get('content-type') || '').startsWith('audio/')) throw new Error(`词库音频样本不可播放:${response.status}`);
  222. }
  223. console.log(JSON.stringify({ completed: true, sourceWords: legacyWords.total, parseWords: finalCount, createdWords: createdCommon, createdAddons, syncedNodes: nodeRequests.length, audioVerified: Boolean(audioSample) }, null, 2));
  224. } finally {
  225. await removeTemporaryAdmin(migrationAdmin);
  226. }