backfill-legacy-users.mjs 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269
  1. #!/usr/bin/env node
  2. import { randomBytes } from 'node:crypto';
  3. import { writeFile } from 'node:fs/promises';
  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. if (!MASTER_KEY) throw new Error('缺少 XIAOSHU_MASTER_KEY');
  11. const headers = {
  12. 'X-Parse-Application-Id': APP_ID,
  13. 'X-Parse-Master-Key': MASTER_KEY,
  14. 'Content-Type': 'application/json',
  15. };
  16. const sleep = (milliseconds) => new Promise((resolve) => setTimeout(resolve, milliseconds));
  17. async function parse(path, init = {}, retries = 4) {
  18. try {
  19. const response = await fetch(`${PARSE_URL}${path}`, { ...init, headers: { ...headers, ...(init.headers || {}) } });
  20. const payload = await response.json().catch(() => ({}));
  21. if (!response.ok || payload.error) {
  22. const error = new Error(typeof payload.error === 'string' ? payload.error : JSON.stringify(payload.error || { status: response.status }));
  23. error.status = response.status;
  24. throw error;
  25. }
  26. return payload;
  27. } catch (error) {
  28. if (!retries || (error.status && error.status < 500 && error.status !== 429)) throw error;
  29. await sleep((5 - retries) * 800);
  30. return parse(path, init, retries - 1);
  31. }
  32. }
  33. const config = (await parse('/config')).params || {};
  34. const legacyUrl = String(config.legacyScheduleApiUrl || '');
  35. const legacyApiId = String(config.legacyScheduleApiId || '');
  36. const legacyApiKey = String(config.legacyScheduleApiKey || '');
  37. if (!legacyUrl || !legacyApiId || !legacyApiKey) throw new Error('生产 Parse Config 未配置旧系统只读接口');
  38. const companyRow = (await parse('/classes/Company?limit=1&keys=objectId')).results?.[0];
  39. if (!companyRow?.objectId) throw new Error('生产 Parse 未找到 Company');
  40. const company = { __type: 'Pointer', className: 'Company', objectId: companyRow.objectId };
  41. async function legacyPage(page) {
  42. const query = new URLSearchParams({ action: 'user_list', uid: '0', psize: '1000', cpage: String(page), apiId: legacyApiId, apiKey: legacyApiKey });
  43. const response = await fetch(`${legacyUrl}${legacyUrl.includes('?') ? '&' : '?'}${query}`, { headers: { accept: 'application/json' } });
  44. const payload = await response.json().catch(() => ({}));
  45. if (!response.ok || Number(payload.retcode) === -1) throw new Error(payload.retmsg || `旧账号接口失败:${response.status}`);
  46. let result = payload.result;
  47. if (typeof result === 'string') result = JSON.parse(result);
  48. return { items: Array.isArray(result) ? result : [], page: payload.page || {} };
  49. }
  50. async function legacyUsers() {
  51. const first = await legacyPage(1);
  52. const total = Number(first.page.itemCount ?? first.items.length);
  53. const pageCount = Math.max(1, Number(first.page.pageCount || Math.ceil(total / 1000)));
  54. const items = [...first.items];
  55. for (let page = 2; page <= pageCount; page += 1) items.push(...(await legacyPage(page)).items);
  56. const byId = new Map();
  57. for (const row of items) {
  58. const id = Number(row.UserID ?? row.userId ?? 0);
  59. if (id > 0) byId.set(id, row);
  60. }
  61. return { total, rows: [...byId.values()], duplicates: items.length - byId.size };
  62. }
  63. async function parseUsers() {
  64. const rows = [];
  65. for (let skip = 0; ; skip += 1000) {
  66. const query = new URLSearchParams({
  67. where: JSON.stringify({ company }), limit: '1000', skip: String(skip), order: 'objectId',
  68. keys: 'objectId,username,legacyUserId,legacyGroupId,legacyUserData,isAdmin,isDeleted,identityType',
  69. });
  70. const page = (await parse(`/users?${query}`)).results || [];
  71. rows.push(...page);
  72. if (page.length < 1000) break;
  73. }
  74. return rows;
  75. }
  76. function text(value) { return String(value ?? '').trim(); }
  77. function number(value, fallback = 0) { const parsed = Number(value); return Number.isFinite(parsed) ? parsed : fallback; }
  78. function first(row, ...keys) { for (const key of keys) if (row[key] !== undefined && row[key] !== null) return row[key]; return undefined; }
  79. function dateIso(value) {
  80. const raw = text(value);
  81. if (!raw) return '';
  82. const normalized = raw.replace(' ', 'T');
  83. const parsed = new Date(normalized + (/Z$|[+-]\d\d:\d\d$/.test(normalized) ? '' : '+08:00'));
  84. return Number.isNaN(parsed.getTime()) ? '' : parsed.toISOString();
  85. }
  86. function safeUsername(raw, legacyId) {
  87. const value = text(raw).slice(0, 120);
  88. return value || `legacy_${legacyId}`;
  89. }
  90. function projection(row) {
  91. const legacyUserId = number(first(row, 'UserID', 'userId'));
  92. const username = safeUsername(first(row, 'UserName', 'username'), legacyUserId);
  93. const nickname = text(first(row, 'HoneyName', 'nickname'));
  94. const realName = text(first(row, 'TrueName', 'RealName', 'realName'));
  95. const mobile = text(first(row, 'Mobile', 'mobile'));
  96. const legacyGroupId = number(first(row, 'GroupID', 'groupId'));
  97. const parentUserId = number(first(row, 'ParentUserID', 'parentUserId'));
  98. const registeredAt = dateIso(first(row, 'RegTime', 'registeredAt'));
  99. const lastLoginAt = dateIso(first(row, 'LastLoginTimes', 'LastLoginTime', 'lastLoginAt'));
  100. const loginCount = number(first(row, 'LoginTimes', 'loginCount'));
  101. const state = number(first(row, 'State', 'state'), 1);
  102. // Deliberate field whitelist: password hashes, payment passwords, security
  103. // questions/answers, sessions and API tokens are never copied.
  104. const legacyUserData = {
  105. UserID: legacyUserId,
  106. UserName: username,
  107. HoneyName: nickname,
  108. TrueName: realName,
  109. Email: text(first(row, 'Email', 'email')),
  110. Mobile: mobile,
  111. GroupID: legacyGroupId,
  112. ParentUserID: parentUserId,
  113. RegTime: registeredAt,
  114. LastLoginTime: lastLoginAt,
  115. LoginTimes: loginCount,
  116. State: state,
  117. Purse: number(first(row, 'Purse', 'purse')),
  118. SilverCoin: number(first(row, 'SilverCoin', 'silverCoin')),
  119. UserExp: number(first(row, 'UserExp', 'userExp')),
  120. UserPoint: number(first(row, 'UserPoint', 'userPoint')),
  121. DummyPurse: number(first(row, 'DummyPurse', 'dummyPurse')),
  122. UserCreit: number(first(row, 'UserCreit', 'Credit', 'credit')),
  123. VIP: number(first(row, 'VIP', 'vip')),
  124. };
  125. const body = {
  126. company, sourceKey: `legacy-sync:user:${legacyUserId}`, legacyUserId, legacyGroupId, legacyUserData,
  127. nickname, realName, mobile, isDisabled: state === 0, isDeleted: false,
  128. };
  129. if (registeredAt) body.registeredAt = { __type: 'Date', iso: registeredAt };
  130. return { legacyUserId, username, body };
  131. }
  132. function requestGroups(requests, size = 40) {
  133. const groups = [];
  134. for (let index = 0; index < requests.length; index += size) groups.push(requests.slice(index, index + size));
  135. return groups;
  136. }
  137. async function executeGroups(groups) {
  138. for (let index = 0; index < groups.length; index += 1) {
  139. const result = await parse('/batch', { method: 'POST', body: JSON.stringify({ requests: groups[index] }) });
  140. const failures = (result || []).filter((item) => item.error);
  141. if (failures.length) throw new Error(`账号批次 ${index + 1} 写入失败:${JSON.stringify(failures.slice(0, 3))}`);
  142. process.stderr.write(`账号写入 ${index + 1}/${groups.length}\n`);
  143. }
  144. }
  145. let tempUserId = '';
  146. let tempFunctionId = '';
  147. let tempToken = '';
  148. let tempPath = '';
  149. async function refreshIdentities() {
  150. const suffix = `${Date.now()}_${randomBytes(4).toString('hex')}`;
  151. const username = `identity_refresh_${suffix}`;
  152. const password = `${randomBytes(24).toString('base64url')}Aa9!`;
  153. tempPath = `xiaoshu/system/refresh-identities-${suffix}`;
  154. const createdUser = await parse('/users', { method: 'POST', body: JSON.stringify({
  155. username, password, isAdmin: true, role: 'admin', roles: ['admin', 'super-admin'], adminRoleKey: 'super-admin',
  156. company, realName: '账号迁移身份刷新临时管理员', testCreatedBy: 'backfill-legacy-users',
  157. }) });
  158. tempUserId = createdUser.objectId;
  159. tempToken = createdUser.sessionToken || (await parse('/login', { method: 'POST', body: JSON.stringify({ username, password }) })).sessionToken;
  160. const code = String.raw`
  161. async function handler(request,response){
  162. try{
  163. const current=request.user||(typeof user!=='undefined'?user:null);if(!current)return response.status(401).json({success:false,message:'需要超级管理员会话'});await current.fetch({useMasterKey:true});
  164. const roles=Array.isArray(current.get('roles'))?current.get('roles').map(String):[];if(current.get('adminRoleKey')!=='super-admin'&&!roles.includes('super-admin'))return response.status(403).json({success:false,message:'仅超级管理员可执行'});
  165. const company=current.get('company'),companyId=company&&company.id;if(!companyId)return response.status(400).json({success:false,message:'管理员缺少帐套'});
  166. const rows=await Psql.query('WITH coach AS (SELECT DISTINCT uid FROM (SELECT NULLIF(CAST("pl" AS text),\'\') uid FROM "CourseAppointment" WHERE "company"=$1 UNION SELECT NULLIF(CAST("fxpl" AS text),\'\') FROM "CourseAppointment" WHERE "company"=$1 UNION SELECT NULLIF(CAST("jsmz" AS text),\'\') FROM "LessonRecord" WHERE "company"=$1) q WHERE uid IS NOT NULL),student AS (SELECT DISTINCT uid FROM (SELECT NULLIF(CAST("szyh" AS text),\'\') uid FROM "CourseAppointment" WHERE "company"=$1 UNION SELECT NULLIF(CAST("userId" AS text),\'\') FROM "DailyStudyRecord" WHERE "company"=$1 UNION SELECT NULLIF(CAST("yhid" AS text),\'\') FROM "CourseBinding" WHERE "company"=$1 UNION SELECT NULLIF(CAST("yhid" AS text),\'\') FROM "PracticeRecord" WHERE "company"=$1) q WHERE uid IS NOT NULL),classified AS (SELECT u."objectId",CASE WHEN COALESCE(CAST(u."legacyGroupId" AS text),u."legacyUserData"->>\'GroupID\')=\'2\' THEN \'store\' WHEN c.uid IS NOT NULL AND s.uid IS NOT NULL THEN \'conflict\' WHEN c.uid IS NOT NULL THEN \'coach\' WHEN s.uid IS NOT NULL OR COALESCE(CAST(u."legacyGroupId" AS text),u."legacyUserData"->>\'GroupID\')=\'1\' THEN \'member\' ELSE \'unknown\' END identity FROM "_User" u LEFT JOIN coach c ON c.uid=COALESCE(CAST(u."legacyUserId" AS text),u."legacyUserData"->>\'UserID\') LEFT JOIN student s ON s.uid=COALESCE(CAST(u."legacyUserId" AS text),u."legacyUserData"->>\'UserID\') WHERE u."company"=$1 AND COALESCE(u."isAdmin",FALSE)=FALSE AND COALESCE(u."isDeleted",FALSE)=FALSE) UPDATE "_User" u SET "identityType"=classified.identity,"updatedAt"=NOW() FROM classified WHERE u."objectId"=classified."objectId" AND u."identityType" IS DISTINCT FROM classified.identity RETURNING classified.identity',[companyId]);
  167. const counts={};for(const row of rows)counts[row.identity]=(counts[row.identity]||0)+1;response.json({success:true,data:{updated:rows.length,counts}});
  168. }catch(error){response.status(Number(error.status)||500).json({success:false,message:String(error.message||error)});}
  169. }`;
  170. const fn = await parse('/classes/Function', { method: 'POST', body: JSON.stringify({
  171. name: tempPath, desc: '一次性刷新迁移账号业务身份', type: 'standalone', path: tempPath, code,
  172. params: [], paramList: [], respType: 'json', respJson: { success: true }, enabled: true,
  173. }) });
  174. tempFunctionId = fn.objectId;
  175. const response = await fetch(`${FUNCTION_URL}/${tempPath}`, {
  176. method: 'POST', headers: { 'X-Parse-Application-Id': APP_ID, 'Content-Type': 'application/json' },
  177. body: JSON.stringify({ token: tempToken, params: {} }),
  178. });
  179. const payload = await response.json().catch(() => ({}));
  180. if (!response.ok || payload.success !== true) throw new Error(payload.message || payload.error || `身份刷新失败:${response.status}`);
  181. return payload.data;
  182. }
  183. async function cleanup() {
  184. if (tempFunctionId) await parse(`/classes/Function/${tempFunctionId}`, { method: 'DELETE' }).catch(() => undefined);
  185. if (tempUserId) {
  186. const where = encodeURIComponent(JSON.stringify({ user: { __type: 'Pointer', className: '_User', objectId: tempUserId } }));
  187. const sessions = await parse(`/classes/_Session?where=${where}&limit=1000&keys=objectId`).catch(() => ({ results: [] }));
  188. for (const session of sessions.results || []) await parse(`/classes/_Session/${session.objectId}`, { method: 'DELETE' }).catch(() => undefined);
  189. await parse(`/users/${tempUserId}`, { method: 'DELETE' }).catch(() => undefined);
  190. }
  191. }
  192. const source = await legacyUsers();
  193. const existing = await parseUsers();
  194. const byLegacyId = new Map();
  195. const byUsername = new Map();
  196. for (const row of existing) {
  197. const legacy = row.legacyUserData || {};
  198. const id = number(row.legacyUserId || legacy.UserID);
  199. if (id > 0) byLegacyId.set(id, row);
  200. if (text(row.username)) byUsername.set(text(row.username).toLowerCase(), row);
  201. }
  202. const requests = [];
  203. const conflicts = [];
  204. let creates = 0;
  205. let attaches = 0;
  206. let updates = 0;
  207. for (const row of source.rows) {
  208. const item = projection(row);
  209. let target = byLegacyId.get(item.legacyUserId);
  210. if (!target) {
  211. const sameUsername = byUsername.get(item.username.toLowerCase());
  212. if (sameUsername && sameUsername.isAdmin !== true && !number(sameUsername.legacyUserId || sameUsername.legacyUserData?.UserID)) {
  213. target = sameUsername;
  214. attaches += 1;
  215. } else if (sameUsername) {
  216. conflicts.push({ legacyUserId: item.legacyUserId, username: item.username, reason: sameUsername.isAdmin === true ? 'username-used-by-admin' : 'username-used-by-another-legacy-id' });
  217. item.username = `legacy_${item.legacyUserId}_${item.username}`.slice(0, 120);
  218. }
  219. }
  220. if (target) {
  221. if (refreshExisting || !number(target.legacyUserId || target.legacyUserData?.UserID)) {
  222. requests.push({ method: 'PUT', path: `/parse/users/${target.objectId}`, body: item.body });
  223. updates += 1;
  224. }
  225. } else {
  226. requests.push({ method: 'POST', path: '/parse/users', body: {
  227. ...item.body, username: item.username, password: `${randomBytes(24).toString('base64url')}Aa9!`,
  228. } });
  229. creates += 1;
  230. }
  231. }
  232. const groups = requestGroups(requests);
  233. const report = {
  234. startedAt: new Date().toISOString(), mode: commit ? 'commit' : 'dry-run', companyId: company.objectId,
  235. legacyTotal: source.total, legacyUnique: source.rows.length, legacyDuplicates: source.duplicates,
  236. parseTotalBefore: existing.length, parseLegacyIdsBefore: byLegacyId.size,
  237. creates, attaches, updates, conflicts, requests: requests.length, batches: groups.length,
  238. };
  239. try {
  240. if (commit && groups.length) await executeGroups(groups);
  241. if (commit) report.identityRefresh = await refreshIdentities();
  242. report.completedAt = new Date().toISOString();
  243. const reportFile = `docs/migration/legacy-user-backfill-${commit ? 'commit' : 'dry-run'}-${new Date().toISOString().replace(/[:.]/g, '-').slice(0, 19)}.json`;
  244. await writeFile(reportFile, `${JSON.stringify(report, null, 2)}\n`, 'utf8');
  245. process.stdout.write(`${JSON.stringify({ ...report, reportFile }, null, 2)}\n`);
  246. } finally {
  247. await cleanup();
  248. }