backfill-core-business-data.mjs 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478
  1. #!/usr/bin/env node
  2. import { writeFile } from 'node:fs/promises';
  3. import { randomBytes } from 'node:crypto';
  4. const APP_ID = process.env.XIAOSHU_PARSE_APP_ID || '7pIbDBJmKx_main';
  5. const MASTER_KEY = process.env.XIAOSHU_MASTER_KEY || '';
  6. const PARSE_URL = (process.env.XIAOSHU_PARSE_URL || 'https://server.xiaoshu.pro/parse').replace(/\/$/, '');
  7. const FUNCTION_URL = PARSE_URL.replace(/\/parse$/, '/api/functions');
  8. const commit = process.argv.includes('--commit');
  9. const refreshExisting = process.argv.includes('--refresh-existing');
  10. const requestedModels = new Set(
  11. (process.argv.find((value) => value.startsWith('--models='))?.split('=')[1] || '53,54,56,58,59,60,61')
  12. .split(',').map(Number).filter(Boolean),
  13. );
  14. const requestedPageStart = Math.max(1, Number(process.argv.find((value) => value.startsWith('--page-start='))?.split('=')[1] || 1));
  15. const requestedPageEndRaw = Number(process.argv.find((value) => value.startsWith('--page-end='))?.split('=')[1] || 0);
  16. const requestedPageEnd = Number.isFinite(requestedPageEndRaw) && requestedPageEndRaw > 0 ? Math.max(requestedPageStart, requestedPageEndRaw) : 0;
  17. const sourceConcurrency = Math.min(10, Math.max(1, Number(process.env.XIAOSHU_SOURCE_CONCURRENCY || 6)));
  18. const writeConcurrency = Math.min(8, Math.max(1, Number(process.env.XIAOSHU_WRITE_CONCURRENCY || 3)));
  19. const pageSize = 1000;
  20. if (!MASTER_KEY) throw new Error('缺少 XIAOSHU_MASTER_KEY');
  21. const specs = {
  22. 53: { label: '练习记录', className: 'PracticeRecord', tableName: 'ZL_C_lxjl', nodeId: 32 },
  23. 54: { label: '预约排课', className: 'CourseAppointment', tableName: 'ZL_C_order', nodeId: 29 },
  24. 56: { label: '每日学习记录', className: 'DailyStudyRecord', tableName: 'ZL_C_ss', nodeId: 291 },
  25. 58: { label: '课程绑定', className: 'CourseBinding', tableName: 'ZL_C_kcbd', nodeId: 28 },
  26. 59: { label: '上课记录', className: 'LessonRecord', tableName: 'ZL_C_skjl', nodeId: 296 },
  27. 60: { label: '抗遗忘记录', className: 'MemoryPracticeRecord', tableName: 'ZL_C_gywjl', nodeId: 388 },
  28. 61: { label: '测评档案', className: 'AssessmentProfile', tableName: 'ZL_C_cpda', nodeId: 389 },
  29. };
  30. for (const modelId of requestedModels) if (!specs[modelId]) throw new Error(`不支持 Model ${modelId}`);
  31. const parseHeaders = {
  32. 'X-Parse-Application-Id': APP_ID,
  33. 'X-Parse-Master-Key': MASTER_KEY,
  34. 'Content-Type': 'application/json',
  35. };
  36. const sleep = (milliseconds) => new Promise((resolve) => setTimeout(resolve, milliseconds));
  37. async function parse(path, init = {}, retries = 5) {
  38. try {
  39. const response = await fetch(`${PARSE_URL}${path}`, {
  40. ...init,
  41. headers: { ...parseHeaders, ...(init.headers || {}) },
  42. });
  43. const payload = await response.json().catch(() => ({}));
  44. if (!response.ok || payload.error) {
  45. const message = typeof payload.error === 'string' ? payload.error : JSON.stringify(payload.error || { status: response.status });
  46. const error = new Error(message);
  47. error.status = response.status;
  48. throw error;
  49. }
  50. return payload;
  51. } catch (error) {
  52. if (!retries || (error.status && error.status < 500 && error.status !== 429)) throw error;
  53. await sleep((6 - retries) * 1000);
  54. return parse(path, init, retries - 1);
  55. }
  56. }
  57. const config = (await parse('/config')).params || {};
  58. const legacyUrl = String(config.legacyScheduleApiUrl || '');
  59. const legacyApiId = String(config.legacyScheduleApiId || '');
  60. const legacyApiKey = String(config.legacyScheduleApiKey || '');
  61. if (!legacyUrl || !legacyApiId || !legacyApiKey) throw new Error('生产 Parse Config 未配置旧系统只读接口');
  62. const companyRow = (await parse('/classes/Company?limit=1&keys=objectId')).results?.[0];
  63. if (!companyRow?.objectId) throw new Error('生产 Parse 未找到 Company');
  64. const company = { __type: 'Pointer', className: 'Company', objectId: companyRow.objectId };
  65. async function legacyPage(modelId, page, retries = 5) {
  66. try {
  67. const query = new URLSearchParams({
  68. action: 'content_list', modelId: String(modelId), psize: String(pageSize), cpage: String(page),
  69. apiId: legacyApiId, apiKey: legacyApiKey,
  70. });
  71. const response = await fetch(`${legacyUrl}${legacyUrl.includes('?') ? '&' : '?'}${query}`, { headers: { accept: 'application/json' } });
  72. const payload = await response.json().catch(() => ({}));
  73. if (!response.ok || Number(payload.retcode) === -1) throw new Error(payload.retmsg || `旧接口失败:${response.status}`);
  74. let result = payload.result;
  75. if (typeof result === 'string') result = JSON.parse(result);
  76. return { items: Array.isArray(result) ? result : [], page: payload.page || {} };
  77. } catch (error) {
  78. if (!retries) throw error;
  79. await sleep((6 - retries) * 800);
  80. return legacyPage(modelId, page, retries - 1);
  81. }
  82. }
  83. async function parallelMap(items, concurrency, worker, progress) {
  84. const output = new Array(items.length);
  85. let cursor = 0;
  86. let completed = 0;
  87. const runners = Array.from({ length: Math.min(concurrency, items.length) }, async () => {
  88. while (true) {
  89. const index = cursor++;
  90. if (index >= items.length) return;
  91. output[index] = await worker(items[index], index);
  92. completed += 1;
  93. if (progress && (completed % progress === 0 || completed === items.length)) process.stderr.write(`读取进度 ${completed}/${items.length}\n`);
  94. }
  95. });
  96. await Promise.all(runners);
  97. return output;
  98. }
  99. async function readLegacyModel(modelId) {
  100. const first = await legacyPage(modelId, 1);
  101. const total = Number(first.page.itemCount ?? first.items.length);
  102. const pageCount = Math.max(1, Number(first.page.pageCount || Math.ceil(total / pageSize)));
  103. const selectedPageStart = Math.min(pageCount, requestedPageStart);
  104. const selectedPageEnd = Math.min(pageCount, requestedPageEnd || pageCount);
  105. const pages = Array.from(
  106. { length: Math.max(0, selectedPageEnd - Math.max(2, selectedPageStart) + 1) },
  107. (_, index) => Math.max(2, selectedPageStart) + index,
  108. );
  109. const rest = await parallelMap(pages, sourceConcurrency, (page) => legacyPage(modelId, page), 50);
  110. const byId = new Map();
  111. const selectedRows = [...(selectedPageStart === 1 ? first.items : []), ...rest.flatMap((page) => page.items)];
  112. for (const row of selectedRows) {
  113. const generalId = Number(row.GeneralID ?? row.generalId ?? 0);
  114. if (generalId) byId.set(generalId, row);
  115. }
  116. return {
  117. total, pageCount, selectedPageStart, selectedPageEnd, selectedCount: selectedRows.length,
  118. rows: [...byId.values()], duplicateOrShifted: selectedRows.length - byId.size,
  119. };
  120. }
  121. async function allParseRows(className, keyField, where, keys) {
  122. const rows = [];
  123. let cursor = -1;
  124. while (true) {
  125. const query = new URLSearchParams({
  126. where: JSON.stringify({ ...where, [keyField]: { $gt: cursor } }),
  127. order: keyField, limit: '1000', keys,
  128. });
  129. const page = (await parse(`/classes/${className}?${query}`)).results || [];
  130. rows.push(...page);
  131. if (page.length < 1000) break;
  132. cursor = Number(page.at(-1)?.[keyField]);
  133. if (!Number.isFinite(cursor)) throw new Error(`${className}.${keyField} 游标无效`);
  134. }
  135. return rows;
  136. }
  137. async function parseRowsByValues(className, keyField, where, values, keys) {
  138. const unique = [...new Set(values.map(Number).filter(Number.isFinite))];
  139. if (!unique.length) return [];
  140. const chunks = [];
  141. for (let index = 0; index < unique.length; index += 400) chunks.push(unique.slice(index, index + 400));
  142. const pages = await parallelMap(chunks, sourceConcurrency, async (chunk) => {
  143. const query = new URLSearchParams({
  144. where: JSON.stringify({ ...where, [keyField]: { $in: chunk } }),
  145. limit: '1000', keys,
  146. });
  147. return (await parse(`/classes/${className}?${query}`)).results || [];
  148. }, 100);
  149. return pages.flat();
  150. }
  151. async function orphanAddonRows(className, keys) {
  152. const rows = [];
  153. for (let skip = 0; ; skip += 1000) {
  154. const query = new URLSearchParams({
  155. where: JSON.stringify({ company, id: { $exists: false } }),
  156. order: 'createdAt', limit: '1000', skip: String(skip), keys,
  157. });
  158. const page = (await parse(`/classes/${className}?${query}`)).results || [];
  159. rows.push(...page);
  160. if (page.length < 1000) break;
  161. }
  162. return rows;
  163. }
  164. function text(value) { return String(value ?? '').trim(); }
  165. function number(value, fallback = 0) { const parsed = Number(value); return Number.isFinite(parsed) ? parsed : fallback; }
  166. function dateValue(value) {
  167. const raw = text(value);
  168. if (!/^\d{4}-\d{1,2}-\d{1,2}(?:[ T]\d{1,2}:\d{2}(?::\d{2}(?:\.\d+)?)?)?/.test(raw)) return null;
  169. const normalized = raw.replace(/^(\d{4})-(\d)(?=-)/, '$1-0$2').replace(/-(\d)(?=[ T])/, '-0$1').replace(' ', 'T');
  170. const date = new Date(normalized + (normalized.includes('T') ? '+08:00' : 'T00:00:00+08:00'));
  171. return Number.isNaN(date.getTime()) ? null : { __type: 'Date', iso: date.toISOString() };
  172. }
  173. function setIf(body, key, value) { if (value !== null && value !== undefined && value !== '') body[key] = value; }
  174. const legacyAddonFields = {
  175. 53: ['kcid','scid','xxcs','yhid','jrscb'],
  176. 54: ['pl','bxrq','dslx','dszt','fxpl','jffs','jssj','kcid','kssj','plxm','scsj','sdsd','szmd','szyh','yysj','yykcid'],
  177. 56: ['pl','con','djq','ygg','dqrq','dsid','fxrl','xxqs','szmdid','userId','learned'],
  178. 58: ['yxx','cksl','kcid','syjd','yhid'],
  179. 59: ['pf','jffs','jsmz','kcid','kclx','kcmc','kzsj','pjnr','pldp','xymz','yyds','plpjsj','szmdid'],
  180. 60: ['fxzt','kcid','plid','wcsj','yhid','kywrq','kywsj','xxjlid','orderId'],
  181. 61: ['df','askid','wrong','userId','answerid','dontKnow','prevScore','totalScore'],
  182. };
  183. const numericAddonFields = new Set(['scid','xxcs','yhid@53','djq','ygg','userId@56','learned','yxx','cksl','pf','kcid@60','plid','yhid@60','xxjlid','orderId']);
  184. function signatureValue(modelId, field, value) {
  185. return numericAddonFields.has(field) || numericAddonFields.has(`${field}@${modelId}`) ? number(value) : text(value);
  186. }
  187. function addonSignature(modelId, source) {
  188. return JSON.stringify(legacyAddonFields[modelId].map((field) => signatureValue(modelId, field, source[field])));
  189. }
  190. function commonBody(modelId, row, itemId) {
  191. const spec = specs[modelId];
  192. const generalId = number(row.GeneralID ?? row.generalId);
  193. const body = {
  194. company,
  195. sourceKey: `legacy:model:${modelId}:general:${generalId}`,
  196. generalId,
  197. itemId,
  198. modelId,
  199. nodeId: number(row.NodeID ?? row.nodeId, spec.nodeId),
  200. tableName: spec.tableName,
  201. title: text(row.Title ?? row.title),
  202. subtitle: text(row.Subtitle ?? row.subtitle),
  203. inputer: text(row.Inputer ?? row.inputer),
  204. topImg: text(row.TopImg ?? row.topImg),
  205. template: text(row.Template ?? row.template),
  206. hits: number(row.Hits ?? row.hits),
  207. status: 99,
  208. };
  209. const created = dateValue(row.CreateTime ?? row.createTime);
  210. if (created) { body.createTime = created; body.upDateTime = created; }
  211. return body;
  212. }
  213. function addonBody(modelId, row, itemId, generalId) {
  214. // Parse REST/SDK reserves the property name `id`. The physical legacy id is
  215. // assigned in one controlled SQL step after both objects exist.
  216. const body = { company, sourceKey: `legacy:model:${modelId}:general:${generalId}` };
  217. if (modelId === 53) Object.assign(body, {
  218. kcid: text(row.kcid), scid: number(row.scid), xxcs: number(row.xxcs), yhid: number(row.yhid), jrscb: text(row.jrscb),
  219. });
  220. if (modelId === 54) for (const field of ['pl','bxrq','dslx','dszt','fxpl','jffs','jssj','kcid','kssj','plxm','scsj','sdsd','szmd','szyh','yysj','yykcid']) body[field] = text(row[field]);
  221. if (modelId === 56) Object.assign(body, {
  222. pl: text(row.pl), con: text(row.con), djq: number(row.djq), ygg: number(row.ygg), dqrq: text(row.dqrq),
  223. dsid: text(row.dsid), fxrl: text(row.fxrl), xxqs: text(row.xxqs), szmdid: text(row.szmdid),
  224. userId: number(row.UserID ?? row.userId), learned: number(row.learned),
  225. });
  226. if (modelId === 58) Object.assign(body, {
  227. yxx: number(row.yxx), cksl: number(row.cksl), kcid: text(row.kcid), syjd: text(row.syjd), yhid: text(row.yhid),
  228. });
  229. if (modelId === 59) {
  230. for (const field of ['jffs','jsmz','kcid','kclx','kcmc','kzsj','pjnr','pldp','xymz','yyds','plpjsj','szmdid']) body[field] = text(row[field]);
  231. setIf(body, 'pf', row.pf === null || row.pf === '' ? null : number(row.pf));
  232. body.legacyGeneralId = generalId;
  233. body.studentName = text(row.Title ?? row.title);
  234. body.coachName = text(row.Inputer ?? row.inputer);
  235. body.courseName = text(row.kcmc ?? row.Subtitle ?? row.subtitle);
  236. body.storeId = number(row.szmdid);
  237. const lessonAt = dateValue(row.plpjsj) || dateValue(row.CreateTime ?? row.createTime);
  238. const sourceCreatedAt = dateValue(row.CreateTime ?? row.createTime);
  239. if (lessonAt) body.lessonAt = lessonAt;
  240. if (sourceCreatedAt) body.sourceCreatedAt = sourceCreatedAt;
  241. }
  242. if (modelId === 60) Object.assign(body, {
  243. fxzt: text(row.fxzt), kcid: number(row.kcid), plid: number(row.plid), wcsj: text(row.wcsj), yhid: number(row.yhid),
  244. kywrq: text(row.kywrq), kywsj: text(row.kywsj), xxjlid: number(row.xxjlid), orderId: number(row.orderId),
  245. });
  246. if (modelId === 61) Object.assign(body, {
  247. df: text(row.df), askid: text(row.askid), wrong: text(row.wrong), userId: text(row.UserID ?? row.userId),
  248. answerid: text(row.answerid), dontKnow: text(row.dontKnow), prevScore: text(row.prev_score ?? row.prevScore), totalScore: text(row.totalScore),
  249. });
  250. return body;
  251. }
  252. function requestGroups(requests, maxRequests = 50, maxBytes = 700_000) {
  253. const groups = [];
  254. let current = [];
  255. let bytes = 0;
  256. for (const request of requests) {
  257. const size = Buffer.byteLength(JSON.stringify(request));
  258. if (current.length && (current.length >= maxRequests || bytes + size > maxBytes)) { groups.push(current); current = []; bytes = 0; }
  259. current.push(request);
  260. bytes += size;
  261. }
  262. if (current.length) groups.push(current);
  263. return groups;
  264. }
  265. async function executeGroups(groups, label) {
  266. let cursor = 0;
  267. let completed = 0;
  268. const runners = Array.from({ length: Math.min(writeConcurrency, groups.length) }, async () => {
  269. while (true) {
  270. const index = cursor++;
  271. if (index >= groups.length) return;
  272. const result = await parse('/batch', { method: 'POST', body: JSON.stringify({ requests: groups[index] }) });
  273. const failures = (result || []).filter((item) => item.error);
  274. if (failures.length) throw new Error(`${label} 批次 ${index + 1} 写入失败:${JSON.stringify(failures.slice(0, 3))}`);
  275. completed += 1;
  276. if (completed % 100 === 0 || completed === groups.length) process.stderr.write(`${label} 写入 ${completed}/${groups.length}\n`);
  277. }
  278. });
  279. await Promise.all(runners);
  280. }
  281. let linkUserId = '';
  282. let linkFunctionId = '';
  283. let linkToken = '';
  284. let linkPath = '';
  285. async function ensureLinkFunction() {
  286. if (linkFunctionId) return;
  287. const suffix = `${Date.now()}_${randomBytes(4).toString('hex')}`;
  288. const username = `backfill_link_${suffix}`;
  289. const password = `${randomBytes(24).toString('base64url')}Aa9!`;
  290. linkPath = `xiaoshu/system/backfill-link-addon-${suffix}`;
  291. const createdUser = await parse('/users', {
  292. method: 'POST',
  293. body: JSON.stringify({
  294. username, password, isAdmin: true, role: 'admin', roles: ['admin', 'super-admin'], adminRoleKey: 'super-admin',
  295. company, realName: '核心业务迁移编号回填临时管理员', testCreatedBy: 'backfill-core-business-data',
  296. }),
  297. });
  298. linkUserId = createdUser.objectId;
  299. linkToken = createdUser.sessionToken || '';
  300. if (!linkToken) linkToken = (await parse('/login', { method: 'POST', body: JSON.stringify({ username, password }) })).sessionToken;
  301. if (!linkToken) throw new Error('临时超级管理员没有会话');
  302. const functionCode = String.raw`
  303. async function handler(request, response) {
  304. try {
  305. const current = request.user || (typeof user !== 'undefined' ? user : null);
  306. if (!current) return response.status(401).json({ success:false, message:'需要超级管理员会话' });
  307. await current.fetch({ useMasterKey:true });
  308. const roles = Array.isArray(current.get('roles')) ? current.get('roles').map(String) : [];
  309. if (current.get('adminRoleKey') !== 'super-admin' && !roles.includes('super-admin')) return response.status(403).json({ success:false, message:'仅超级管理员可执行' });
  310. const modelId = Number(request.params && request.params.modelId || request.body && request.body.params && request.body.params.modelId || 0);
  311. const tableMap = ${JSON.stringify(Object.fromEntries(Object.entries(specs).map(([modelId, spec]) => [modelId, spec.className])))};
  312. const table = tableMap[String(modelId)];
  313. if (!table) return response.status(400).json({ success:false, message:'不支持该业务模型' });
  314. const company = current.get('company');
  315. const companyId = company && company.id;
  316. if (!companyId) return response.status(400).json({ success:false, message:'管理员缺少帐套' });
  317. const rows = await Psql.query(
  318. 'UPDATE "' + table + '" a SET "id"=c."itemId","updatedAt"=NOW() FROM "CommonModel" c WHERE a."company"=$1 AND c."company"=$1 AND CAST(c."modelId" AS text)=$2 AND (a."sourceKey"=c."sourceKey" OR a."sourceKey"=(\'legacy:model:\'||$2||\':general:\'||CAST(c."generalId" AS text))) AND a."id" IS DISTINCT FROM c."itemId" RETURNING a."objectId"',
  319. [companyId, String(modelId)]
  320. );
  321. response.json({ success:true, data:{ modelId, table, linked:rows.length } });
  322. } catch (error) {
  323. response.status(Number(error.status)||500).json({ success:false, message:String(error.message||error) });
  324. }
  325. }`;
  326. const createdFunction = await parse('/classes/Function', {
  327. method: 'POST',
  328. body: JSON.stringify({
  329. name: linkPath, desc: '一次性回填核心业务附表旧编号', type: 'standalone', path: linkPath, code: functionCode,
  330. params: [], paramList: [], respType: 'json', respJson: { success: true }, enabled: true,
  331. }),
  332. });
  333. linkFunctionId = createdFunction.objectId;
  334. }
  335. async function linkAddonIds(modelId) {
  336. await ensureLinkFunction();
  337. const response = await fetch(`${FUNCTION_URL}/${linkPath}`, {
  338. method: 'POST',
  339. headers: { 'X-Parse-Application-Id': APP_ID, 'Content-Type': 'application/json' },
  340. body: JSON.stringify({ token: linkToken, params: { modelId } }),
  341. });
  342. const payload = await response.json().catch(() => ({}));
  343. if (!response.ok || payload.success !== true) throw new Error(payload.message || payload.error || `附表编号回填失败:${response.status}`);
  344. return Number(payload.data?.linked || 0);
  345. }
  346. async function cleanupLinkFunction() {
  347. if (linkFunctionId) await parse(`/classes/Function/${linkFunctionId}`, { method: 'DELETE' }).catch(() => undefined);
  348. if (linkUserId) {
  349. const where = encodeURIComponent(JSON.stringify({ user: { __type: 'Pointer', className: '_User', objectId: linkUserId } }));
  350. const sessions = await parse(`/classes/_Session?where=${where}&limit=1000&keys=objectId`).catch(() => ({ results: [] }));
  351. for (const session of sessions.results || []) await parse(`/classes/_Session/${session.objectId}`, { method: 'DELETE' }).catch(() => undefined);
  352. await parse(`/users/${linkUserId}`, { method: 'DELETE' }).catch(() => undefined);
  353. }
  354. }
  355. const report = { startedAt: new Date().toISOString(), mode: commit ? 'commit' : 'dry-run', companyId: company.objectId, models: [] };
  356. try {
  357. for (const modelId of [...requestedModels].sort((a, b) => a - b)) {
  358. const spec = specs[modelId];
  359. process.stderr.write(`\n读取旧系统 ${spec.label}(Model ${modelId})…\n`);
  360. const source = await readLegacyModel(modelId);
  361. process.stderr.write(`旧系统分片 ${source.rows.length} 条(总计 ${source.total}),页 ${source.selectedPageStart}-${source.selectedPageEnd}/${source.pageCount}\n`);
  362. const sliced = source.selectedPageStart > 1 || source.selectedPageEnd < source.pageCount;
  363. const sourceGeneralIds = source.rows.map((row) => number(row.GeneralID ?? row.generalId));
  364. const commonRows = sliced
  365. ? await parseRowsByValues('CommonModel', 'generalId', { company, modelId }, sourceGeneralIds, 'objectId,generalId,itemId,sourceKey')
  366. : await allParseRows('CommonModel', 'generalId', { company, modelId }, 'objectId,generalId,itemId,sourceKey');
  367. const commonByGeneralIdForCandidates = new Map(commonRows.map((row) => [number(row.generalId), row]));
  368. const addonCandidateIds = sourceGeneralIds.map((generalId) => number(commonByGeneralIdForCandidates.get(generalId)?.itemId, generalId));
  369. const addonRows = sliced
  370. ? await parseRowsByValues(spec.className, 'id', { company }, addonCandidateIds, 'objectId,id,sourceKey')
  371. : await allParseRows(spec.className, 'id', { company }, 'objectId,id,sourceKey');
  372. const commonByGeneralId = new Map(commonRows.map((row) => [number(row.generalId), row]));
  373. const addonById = new Map(addonRows.map((row) => [number(row.id), row]));
  374. const usedAddonIds = new Set(addonById.keys());
  375. const requests = [];
  376. const orphanRows = await orphanAddonRows(spec.className, ['objectId', ...legacyAddonFields[modelId]].join(','));
  377. const existingCandidatesBySignature = new Map();
  378. for (const row of source.rows) {
  379. const generalId = number(row.GeneralID ?? row.generalId);
  380. const common = commonByGeneralId.get(generalId);
  381. if (!common) continue;
  382. const signature = addonSignature(modelId, addonBody(modelId, row, number(common.itemId), generalId));
  383. if (!existingCandidatesBySignature.has(signature)) existingCandidatesBySignature.set(signature, []);
  384. existingCandidatesBySignature.get(signature).push({ generalId, itemId: number(common.itemId), source: row });
  385. }
  386. for (const candidates of existingCandidatesBySignature.values()) candidates.sort((a, b) => a.itemId - b.itemId);
  387. let repairedOrphans = 0;
  388. let unmatchedOrphans = 0;
  389. for (const orphan of orphanRows) {
  390. const candidates = existingCandidatesBySignature.get(addonSignature(modelId, orphan)) || [];
  391. const candidate = candidates.find((item) => !addonById.has(item.itemId));
  392. if (!candidate) { unmatchedOrphans += 1; continue; }
  393. const body = addonBody(modelId, candidate.source, candidate.itemId, candidate.generalId);
  394. requests.push({ method: 'PUT', path: `/parse/classes/${spec.className}/${orphan.objectId}`, body });
  395. addonById.set(candidate.itemId, { ...orphan, ...body });
  396. usedAddonIds.add(candidate.itemId);
  397. repairedOrphans += 1;
  398. }
  399. let createCommon = 0;
  400. let createAddon = 0;
  401. let updateCommon = 0;
  402. let updateAddon = 0;
  403. for (const row of source.rows) {
  404. const generalId = number(row.GeneralID ?? row.generalId);
  405. const existingCommon = commonByGeneralId.get(generalId);
  406. let itemId = existingCommon ? number(existingCommon.itemId) : generalId;
  407. if (!existingCommon && usedAddonIds.has(itemId)) itemId = 900_000_000_000 + generalId * 100 + modelId;
  408. usedAddonIds.add(itemId);
  409. const existingAddon = addonById.get(itemId);
  410. if (!existingAddon) {
  411. requests.push({ method: 'POST', path: `/parse/classes/${spec.className}`, body: addonBody(modelId, row, itemId, generalId) });
  412. createAddon += 1;
  413. } else if (refreshExisting) {
  414. requests.push({ method: 'PUT', path: `/parse/classes/${spec.className}/${existingAddon.objectId}`, body: addonBody(modelId, row, itemId, generalId) });
  415. updateAddon += 1;
  416. }
  417. if (!existingCommon) {
  418. requests.push({ method: 'POST', path: '/parse/classes/CommonModel', body: commonBody(modelId, row, itemId) });
  419. createCommon += 1;
  420. } else if (refreshExisting) {
  421. requests.push({ method: 'PUT', path: `/parse/classes/CommonModel/${existingCommon.objectId}`, body: commonBody(modelId, row, itemId) });
  422. updateCommon += 1;
  423. }
  424. }
  425. const groups = requestGroups(requests);
  426. const summary = {
  427. modelId, label: spec.label, sourceTotal: source.total, sourceUnique: source.rows.length,
  428. sourcePageStart: source.selectedPageStart, sourcePageEnd: source.selectedPageEnd, sourcePageCount: source.pageCount,
  429. sourceDuplicateOrShifted: source.duplicateOrShifted, parseCommonBefore: commonRows.length, parseAddonBefore: addonRows.length,
  430. orphanRows: orphanRows.length, repairedOrphans, unmatchedOrphans,
  431. createCommon, createAddon, updateCommon, updateAddon, requests: requests.length, batches: groups.length,
  432. };
  433. report.models.push(summary);
  434. process.stderr.write(`${JSON.stringify(summary)}\n`);
  435. if (commit && groups.length) {
  436. await executeGroups(groups, spec.label);
  437. summary.linkedAddonIds = await linkAddonIds(modelId);
  438. process.stderr.write(`${spec.label} 附表编号回填 ${summary.linkedAddonIds} 条\n`);
  439. }
  440. }
  441. report.completedAt = new Date().toISOString();
  442. const reportFile = `docs/migration/core-business-backfill-${commit ? 'commit' : 'dry-run'}-${new Date().toISOString().replace(/[:.]/g, '-').slice(0, 19)}.json`;
  443. await writeFile(reportFile, `${JSON.stringify(report, null, 2)}\n`, 'utf8');
  444. process.stdout.write(`${JSON.stringify({ ...report, reportFile }, null, 2)}\n`);
  445. } finally {
  446. await cleanupLinkFunction();
  447. }