merge-reading-users.mjs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256
  1. #!/usr/bin/env node
  2. import { randomBytes } from 'node:crypto';
  3. const APP_ID = process.env.XIAOSHU_PARSE_APP_ID || '7pIbDBJmKx_main';
  4. const MASTER_KEY = process.env.XIAOSHU_MASTER_KEY || '';
  5. const PARSE_URL = (process.env.XIAOSHU_PARSE_URL || 'https://server.xiaoshu.pro/parse').replace(/\/$/, '');
  6. const FUNCTION_URL = PARSE_URL.replace(/\/parse$/, '/api/functions');
  7. const COMPANY_ID = process.env.XIAOSHU_COMPANY_ID || '7pIbDBJmKx';
  8. const commit = process.argv.includes('--commit');
  9. if (!MASTER_KEY) throw new Error('缺少 XIAOSHU_MASTER_KEY');
  10. const headers = {
  11. 'X-Parse-Application-Id': APP_ID,
  12. 'X-Parse-Master-Key': MASTER_KEY,
  13. 'Content-Type': 'application/json',
  14. };
  15. const company = { __type: 'Pointer', className: 'Company', objectId: COMPANY_ID };
  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}`, {
  20. ...init,
  21. headers: { ...headers, ...(init.headers || {}) },
  22. });
  23. const payload = await response.json().catch(() => ({}));
  24. if (!response.ok || payload.error) {
  25. const error = new Error(typeof payload.error === 'string' ? payload.error : JSON.stringify(payload.error || { status: response.status }));
  26. error.status = response.status;
  27. throw error;
  28. }
  29. return payload;
  30. } catch (error) {
  31. if (!retries || (error.status && error.status < 500 && error.status !== 429)) throw error;
  32. await sleep((5 - retries) * 800);
  33. return parse(path, init, retries - 1);
  34. }
  35. }
  36. async function allRows(path, where, keys) {
  37. const rows = [];
  38. for (let skip = 0; ; skip += 1000) {
  39. const query = new URLSearchParams({
  40. where: JSON.stringify(where),
  41. limit: '1000',
  42. skip: String(skip),
  43. order: 'objectId',
  44. keys,
  45. });
  46. const page = (await parse(`${path}?${query}`)).results || [];
  47. rows.push(...page);
  48. if (page.length < 1000) break;
  49. }
  50. return rows;
  51. }
  52. function text(value) { return String(value ?? '').trim(); }
  53. function number(value) { const parsed = Number(value); return Number.isFinite(parsed) ? parsed : 0; }
  54. function legacyId(row) { return number(row.legacyUserId || row.legacyUserData?.UserID); }
  55. function normalizePhone(value) {
  56. const digits = text(value).replace(/\D/g, '');
  57. if (/^86\d{11}$/.test(digits)) return digits.slice(2);
  58. return /^1\d{10}$/.test(digits) ? digits : '';
  59. }
  60. function phones(row) {
  61. return new Set([
  62. row.mobile, row.phone, row.mobilePhoneNumber, row.username,
  63. row.legacyUserData?.Mobile, row.legacyUserData?.UserName,
  64. ].map(normalizePhone).filter(Boolean));
  65. }
  66. function normalizeName(value) {
  67. return text(value).toLowerCase().replace(/[\s\p{P}\p{S}]+/gu, '');
  68. }
  69. function names(row) {
  70. return new Set([
  71. row.nickname, row.nickName, row.realName, row.name,
  72. row.legacyUserData?.HoneyName, row.legacyUserData?.TrueName,
  73. ].map(normalizeName).filter((value) => value.length >= 2));
  74. }
  75. function intersects(left, right) { return [...left].some((value) => right.has(value)); }
  76. function isCanonical(row) { return text(row.sourceKey).startsWith('legacy-sync:user:') && !row.isDeleted; }
  77. function isReadingImport(row) {
  78. return !text(row.sourceKey) && !legacyId(row) && row.isAdmin !== true && !row.isDeleted;
  79. }
  80. function buildPlan(users) {
  81. const targets = users.filter(isCanonical);
  82. const sources = users.filter(isReadingImport);
  83. const targetsByPhone = new Map();
  84. for (const target of targets) {
  85. for (const phone of phones(target)) {
  86. const values = targetsByPhone.get(phone) || [];
  87. values.push(target);
  88. targetsByPhone.set(phone, values);
  89. }
  90. }
  91. const merge = [];
  92. const promote = [];
  93. let nameResolved = 0;
  94. let ambiguous = 0;
  95. let unmatched = 0;
  96. for (const source of sources) {
  97. const candidates = new Map();
  98. for (const phone of phones(source)) {
  99. for (const candidate of targetsByPhone.get(phone) || []) candidates.set(candidate.objectId, candidate);
  100. }
  101. let possible = [...candidates.values()];
  102. if (possible.length > 1) {
  103. const sourceNames = names(source);
  104. const nameMatches = possible.filter((candidate) => intersects(sourceNames, names(candidate)));
  105. if (nameMatches.length === 1) { possible = nameMatches; nameResolved += 1; }
  106. }
  107. if (possible.length === 1) {
  108. merge.push({ source: source.objectId, target: possible[0].objectId });
  109. } else {
  110. promote.push(source.objectId);
  111. if (possible.length > 1) ambiguous += 1;
  112. else unmatched += 1;
  113. }
  114. }
  115. return { targets, sources, merge, promote, nameResolved, ambiguous, unmatched };
  116. }
  117. async function ownerStats() {
  118. const result = {};
  119. for (const className of ['SurveyItem', 'SurveyLog']) {
  120. const rows = await allRows(`/classes/${className}`, { user: { $exists: true } }, 'objectId,user');
  121. const counts = new Map();
  122. for (const row of rows) {
  123. const id = text(row.user?.objectId || row.user).replace(/^_User\$/, '');
  124. if (id) counts.set(id, (counts.get(id) || 0) + 1);
  125. }
  126. result[className] = { records: rows.length, owners: counts.size, counts };
  127. }
  128. return result;
  129. }
  130. let temporaryUserId = '';
  131. let temporaryFunctionId = '';
  132. async function executePlan(plan) {
  133. const suffix = `${Date.now()}_${randomBytes(4).toString('hex')}`;
  134. const username = `reading_merge_${suffix}`;
  135. const password = `${randomBytes(24).toString('base64url')}Aa9!`;
  136. const path = `xiaoshu/system/merge-reading-users-${suffix}`;
  137. const created = await parse('/users', {
  138. method: 'POST',
  139. body: JSON.stringify({
  140. username, password, company, isAdmin: true, role: 'admin',
  141. roles: ['admin', 'super-admin'], adminRoleKey: 'super-admin',
  142. realName: '阅读账号合并临时管理员', testCreatedBy: 'merge-reading-users',
  143. }),
  144. });
  145. temporaryUserId = created.objectId;
  146. const token = created.sessionToken || (await parse('/login', {
  147. method: 'POST', body: JSON.stringify({ username, password }),
  148. })).sessionToken;
  149. const embeddedMerge = JSON.stringify(plan.merge).replace(/</g, '\\u003c');
  150. const embeddedPromote = JSON.stringify(plan.promote).replace(/</g, '\\u003c');
  151. const code = String.raw`
  152. async function handler(request,response){
  153. try{
  154. const current=request.user||(typeof user!=='undefined'?user:null);if(!current)return response.status(401).json({success:false,message:'需要超级管理员会话'});await current.fetch({useMasterKey:true});
  155. 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:'仅超级管理员可执行'});
  156. const companyId=${JSON.stringify(COMPANY_ID)},merge=${embeddedMerge},promote=${embeddedPromote};
  157. await Psql.none('ALTER TABLE "SurveyLog" ALTER COLUMN "user" TYPE text USING CAST("user" AS text)');
  158. const merged=merge.length?await Psql.one('WITH input AS MATERIALIZED (SELECT x.source,x.target FROM jsonb_to_recordset($2::jsonb) AS x(source text,target text)),valid AS MATERIALIZED (SELECT i.source,i.target FROM input i JOIN "_User" s ON s."objectId"=i.source AND s."company"=$1 JOIN "_User" t ON t."objectId"=i.target AND t."company"=$1 WHERE COALESCE(s."isDeleted",FALSE)=FALSE AND COALESCE(t."isDeleted",FALSE)=FALSE),items AS (UPDATE "SurveyItem" r SET "user"=v.target,"updatedAt"=NOW() FROM valid v WHERE CAST(r."user" AS text) IN (v.source,\'_User$\'||v.source) RETURNING r."objectId"),logs AS (UPDATE "SurveyLog" r SET "user"=v.target,"updatedAt"=NOW() FROM valid v WHERE CAST(r."user" AS text) IN (v.source,\'_User$\'||v.source) RETURNING r."objectId"),profiles AS (UPDATE "_User" t SET "level"=COALESCE(NULLIF(t."level",\'\'),NULLIF(s."level",\'\')),"updatedAt"=NOW() FROM valid v JOIN "_User" s ON s."objectId"=v.source WHERE t."objectId"=v.target RETURNING t."objectId"),sessions AS (DELETE FROM "_Session" x USING valid v WHERE CAST(x."user" AS text) IN (v.source,\'_User$\'||v.source) RETURNING x."objectId"),sources AS (UPDATE "_User" s SET "isDisabled"=TRUE,"isDeleted"=TRUE,"identityType"=\'merged\',"sourceKey"=\'reading-merged:user:\'||s."objectId","legacyUserData"=COALESCE(s."legacyUserData",\'{}\'::jsonb)||jsonb_build_object(\'ReadingMergedIntoObjectId\',v.target,\'ReadingMergedAt\',NOW()),"updatedAt"=NOW() FROM valid v WHERE s."objectId"=v.source RETURNING s."objectId") SELECT (SELECT COUNT(*)::int FROM sources) users,(SELECT COUNT(*)::int FROM items) items,(SELECT COUNT(*)::int FROM logs) logs,(SELECT COUNT(*)::int FROM sessions) sessions',[companyId,JSON.stringify(merge)]):{users:0,items:0,logs:0,sessions:0};
  159. const promoted=promote.length?await Psql.one('WITH lock_row AS MATERIALIZED (SELECT pg_advisory_xact_lock(hashtext(\'xiaoshu-legacy-user-id\'))),input AS MATERIALIZED (SELECT value#>>\'{}\' object_id FROM jsonb_array_elements($2::jsonb)),maximum AS MATERIALIZED (SELECT COALESCE(MAX(CASE WHEN COALESCE("legacyUserId",0)>0 THEN "legacyUserId" ELSE 0 END),0)::bigint max_id FROM "_User",lock_row WHERE "company"=$1),candidates AS MATERIALIZED (SELECT u."objectId",(SELECT max_id FROM maximum)+ROW_NUMBER() OVER(ORDER BY u."objectId") legacy_id FROM "_User" u JOIN input i ON i.object_id=u."objectId" WHERE u."company"=$1 AND COALESCE(u."isDeleted",FALSE)=FALSE AND COALESCE(u."legacyUserId",0)=0),updated AS (UPDATE "_User" u SET "legacyUserId"=c.legacy_id,"legacyGroupId"=1,"identityType"=\'member\',"sourceKey"=\'reading-unified:user:\'||u."objectId","isDisabled"=FALSE,"isDeleted"=FALSE,"registeredAt"=COALESCE(u."registeredAt",u."createdAt"),"legacyUserData"=COALESCE(u."legacyUserData",\'{}\'::jsonb)||jsonb_build_object(\'UserID\',c.legacy_id,\'UserName\',COALESCE(NULLIF(u."username",\'\'),\'reading_\'||u."objectId"),\'HoneyName\',COALESCE(NULLIF(u."nickname",\'\'),NULLIF(u."realName",\'\'),NULLIF(u."username",\'\'),\'阅读会员\'),\'Mobile\',COALESCE(NULLIF(u."mobile",\'\'),NULLIF(u."phone",\'\'),CASE WHEN u."username"~\'^1[0-9]{10}$\' THEN u."username" ELSE \'\' END),\'GroupID\',1,\'ParentUserID\',0,\'RegTime\',COALESCE(u."registeredAt",u."createdAt"),\'State\',1,\'ReadingUnifiedAt\',NOW()),"updatedAt"=NOW() FROM candidates c WHERE u."objectId"=c."objectId" RETURNING u."objectId") SELECT COUNT(*)::int users FROM updated',[companyId,JSON.stringify(promote)]):{users:0};
  160. response.json({success:true,data:{merged,promoted}});
  161. }catch(error){response.status(Number(error.status)||500).json({success:false,message:String(error.message||error)});}
  162. }`;
  163. const fn = await parse('/classes/Function', {
  164. method: 'POST',
  165. body: JSON.stringify({
  166. name: path, desc: '一次性合并阅读账号与原会员', type: 'standalone', path, code,
  167. params: [], paramList: [], respType: 'json', respJson: { success: true }, enabled: true,
  168. }),
  169. });
  170. temporaryFunctionId = fn.objectId;
  171. const response = await fetch(`${FUNCTION_URL}/${path}`, {
  172. method: 'POST',
  173. headers: { 'X-Parse-Application-Id': APP_ID, 'Content-Type': 'application/json' },
  174. body: JSON.stringify({ token, params: {} }),
  175. });
  176. const payload = await response.json().catch(() => ({}));
  177. if (!response.ok || payload.success !== true) throw new Error(payload.message || payload.error || `账号合并失败:${response.status}`);
  178. return payload.data;
  179. }
  180. async function cleanup() {
  181. if (temporaryFunctionId) await parse(`/classes/Function/${temporaryFunctionId}`, { method: 'DELETE' }).catch(() => undefined);
  182. if (temporaryUserId) {
  183. const where = encodeURIComponent(JSON.stringify({ user: { __type: 'Pointer', className: '_User', objectId: temporaryUserId } }));
  184. const sessions = await parse(`/classes/_Session?where=${where}&limit=1000&keys=objectId`).catch(() => ({ results: [] }));
  185. for (const session of sessions.results || []) await parse(`/classes/_Session/${session.objectId}`, { method: 'DELETE' }).catch(() => undefined);
  186. await parse(`/users/${temporaryUserId}`, { method: 'DELETE' }).catch(() => undefined);
  187. }
  188. }
  189. const userKeys = 'objectId,username,nickname,nickName,realName,name,mobile,phone,mobilePhoneNumber,level,company,sourceKey,legacyUserId,legacyGroupId,legacyUserData,identityType,isAdmin,isDisabled,isDeleted,createdAt,registeredAt';
  190. const usersBefore = await allRows('/users', { company }, userKeys);
  191. const ownersBefore = await ownerStats();
  192. const plan = buildPlan(usersBefore);
  193. const sourceIds = new Set(plan.sources.map((row) => row.objectId));
  194. const readingSources = new Set();
  195. for (const value of Object.values(ownersBefore)) {
  196. for (const owner of value.counts.keys()) if (sourceIds.has(owner)) readingSources.add(owner);
  197. }
  198. const summary = {
  199. mode: commit ? 'commit' : 'dry-run',
  200. usersBefore: usersBefore.length,
  201. canonicalUsersBefore: plan.targets.length,
  202. readingImports: plan.sources.length,
  203. readingImportsWithRecords: readingSources.size,
  204. mergeUsers: plan.merge.length,
  205. promoteUsers: plan.promote.length,
  206. maxMergeSourceIdLength: Math.max(0, ...plan.merge.map((item) => item.source.length)),
  207. maxMergeTargetIdLength: Math.max(0, ...plan.merge.map((item) => item.target.length)),
  208. nameResolved: plan.nameResolved,
  209. ambiguousPromoted: plan.ambiguous,
  210. unmatchedPromoted: plan.unmatched,
  211. surveyItemsToMove: plan.merge.reduce((total, item) => total + (ownersBefore.SurveyItem.counts.get(item.source) || 0), 0),
  212. surveyLogsToMove: plan.merge.reduce((total, item) => total + (ownersBefore.SurveyLog.counts.get(item.source) || 0), 0),
  213. };
  214. try {
  215. if (commit && (plan.merge.length || plan.promote.length)) summary.executed = await executePlan(plan);
  216. if (commit) {
  217. const usersAfter = await allRows('/users', { company }, userKeys);
  218. const ownersAfter = await ownerStats();
  219. const remaining = usersAfter.filter(isReadingImport);
  220. const activeById = new Map(usersAfter.filter((row) => !row.isDeleted).map((row) => [row.objectId, row]));
  221. const invalidOwners = new Set();
  222. for (const value of Object.values(ownersAfter)) {
  223. for (const owner of value.counts.keys()) {
  224. const row = activeById.get(owner);
  225. if (!row || (!isCanonical(row) && !text(row.sourceKey).startsWith('reading-unified:user:'))) invalidOwners.add(owner);
  226. }
  227. }
  228. summary.verification = {
  229. usersAfter: usersAfter.length,
  230. remainingUnnormalizedUsers: remaining.length,
  231. invalidReadingOwners: invalidOwners.size,
  232. surveyItems: ownersAfter.SurveyItem.records,
  233. surveyLogs: ownersAfter.SurveyLog.records,
  234. };
  235. if (remaining.length || invalidOwners.size) throw new Error(`合并后验证失败:未归一账号 ${remaining.length},无效阅读数据归属 ${invalidOwners.size}`);
  236. }
  237. process.stdout.write(`${JSON.stringify(summary, null, 2)}\n`);
  238. } finally {
  239. await cleanup();
  240. }