| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256 |
- #!/usr/bin/env node
- import { randomBytes } from 'node:crypto';
- const APP_ID = process.env.XIAOSHU_PARSE_APP_ID || '7pIbDBJmKx_main';
- const MASTER_KEY = process.env.XIAOSHU_MASTER_KEY || '';
- const PARSE_URL = (process.env.XIAOSHU_PARSE_URL || 'https://server.xiaoshu.pro/parse').replace(/\/$/, '');
- const FUNCTION_URL = PARSE_URL.replace(/\/parse$/, '/api/functions');
- const COMPANY_ID = process.env.XIAOSHU_COMPANY_ID || '7pIbDBJmKx';
- const commit = process.argv.includes('--commit');
- if (!MASTER_KEY) throw new Error('缺少 XIAOSHU_MASTER_KEY');
- const headers = {
- 'X-Parse-Application-Id': APP_ID,
- 'X-Parse-Master-Key': MASTER_KEY,
- 'Content-Type': 'application/json',
- };
- const company = { __type: 'Pointer', className: 'Company', objectId: COMPANY_ID };
- const sleep = (milliseconds) => new Promise((resolve) => setTimeout(resolve, milliseconds));
- async function parse(path, init = {}, retries = 4) {
- try {
- const response = await fetch(`${PARSE_URL}${path}`, {
- ...init,
- headers: { ...headers, ...(init.headers || {}) },
- });
- const payload = await response.json().catch(() => ({}));
- if (!response.ok || payload.error) {
- const error = new Error(typeof payload.error === 'string' ? payload.error : JSON.stringify(payload.error || { status: response.status }));
- error.status = response.status;
- throw error;
- }
- return payload;
- } catch (error) {
- if (!retries || (error.status && error.status < 500 && error.status !== 429)) throw error;
- await sleep((5 - retries) * 800);
- return parse(path, init, retries - 1);
- }
- }
- async function allRows(path, where, keys) {
- const rows = [];
- for (let skip = 0; ; skip += 1000) {
- const query = new URLSearchParams({
- where: JSON.stringify(where),
- limit: '1000',
- skip: String(skip),
- order: 'objectId',
- keys,
- });
- const page = (await parse(`${path}?${query}`)).results || [];
- rows.push(...page);
- if (page.length < 1000) break;
- }
- return rows;
- }
- function text(value) { return String(value ?? '').trim(); }
- function number(value) { const parsed = Number(value); return Number.isFinite(parsed) ? parsed : 0; }
- function legacyId(row) { return number(row.legacyUserId || row.legacyUserData?.UserID); }
- function normalizePhone(value) {
- const digits = text(value).replace(/\D/g, '');
- if (/^86\d{11}$/.test(digits)) return digits.slice(2);
- return /^1\d{10}$/.test(digits) ? digits : '';
- }
- function phones(row) {
- return new Set([
- row.mobile, row.phone, row.mobilePhoneNumber, row.username,
- row.legacyUserData?.Mobile, row.legacyUserData?.UserName,
- ].map(normalizePhone).filter(Boolean));
- }
- function normalizeName(value) {
- return text(value).toLowerCase().replace(/[\s\p{P}\p{S}]+/gu, '');
- }
- function names(row) {
- return new Set([
- row.nickname, row.nickName, row.realName, row.name,
- row.legacyUserData?.HoneyName, row.legacyUserData?.TrueName,
- ].map(normalizeName).filter((value) => value.length >= 2));
- }
- function intersects(left, right) { return [...left].some((value) => right.has(value)); }
- function isCanonical(row) { return text(row.sourceKey).startsWith('legacy-sync:user:') && !row.isDeleted; }
- function isReadingImport(row) {
- return !text(row.sourceKey) && !legacyId(row) && row.isAdmin !== true && !row.isDeleted;
- }
- function buildPlan(users) {
- const targets = users.filter(isCanonical);
- const sources = users.filter(isReadingImport);
- const targetsByPhone = new Map();
- for (const target of targets) {
- for (const phone of phones(target)) {
- const values = targetsByPhone.get(phone) || [];
- values.push(target);
- targetsByPhone.set(phone, values);
- }
- }
- const merge = [];
- const promote = [];
- let nameResolved = 0;
- let ambiguous = 0;
- let unmatched = 0;
- for (const source of sources) {
- const candidates = new Map();
- for (const phone of phones(source)) {
- for (const candidate of targetsByPhone.get(phone) || []) candidates.set(candidate.objectId, candidate);
- }
- let possible = [...candidates.values()];
- if (possible.length > 1) {
- const sourceNames = names(source);
- const nameMatches = possible.filter((candidate) => intersects(sourceNames, names(candidate)));
- if (nameMatches.length === 1) { possible = nameMatches; nameResolved += 1; }
- }
- if (possible.length === 1) {
- merge.push({ source: source.objectId, target: possible[0].objectId });
- } else {
- promote.push(source.objectId);
- if (possible.length > 1) ambiguous += 1;
- else unmatched += 1;
- }
- }
- return { targets, sources, merge, promote, nameResolved, ambiguous, unmatched };
- }
- async function ownerStats() {
- const result = {};
- for (const className of ['SurveyItem', 'SurveyLog']) {
- const rows = await allRows(`/classes/${className}`, { user: { $exists: true } }, 'objectId,user');
- const counts = new Map();
- for (const row of rows) {
- const id = text(row.user?.objectId || row.user).replace(/^_User\$/, '');
- if (id) counts.set(id, (counts.get(id) || 0) + 1);
- }
- result[className] = { records: rows.length, owners: counts.size, counts };
- }
- return result;
- }
- let temporaryUserId = '';
- let temporaryFunctionId = '';
- async function executePlan(plan) {
- const suffix = `${Date.now()}_${randomBytes(4).toString('hex')}`;
- const username = `reading_merge_${suffix}`;
- const password = `${randomBytes(24).toString('base64url')}Aa9!`;
- const path = `xiaoshu/system/merge-reading-users-${suffix}`;
- const created = await parse('/users', {
- method: 'POST',
- body: JSON.stringify({
- username, password, company, isAdmin: true, role: 'admin',
- roles: ['admin', 'super-admin'], adminRoleKey: 'super-admin',
- realName: '阅读账号合并临时管理员', testCreatedBy: 'merge-reading-users',
- }),
- });
- temporaryUserId = created.objectId;
- const token = created.sessionToken || (await parse('/login', {
- method: 'POST', body: JSON.stringify({ username, password }),
- })).sessionToken;
- const embeddedMerge = JSON.stringify(plan.merge).replace(/</g, '\\u003c');
- const embeddedPromote = JSON.stringify(plan.promote).replace(/</g, '\\u003c');
- const code = String.raw`
- async function handler(request,response){
- try{
- const current=request.user||(typeof user!=='undefined'?user:null);if(!current)return response.status(401).json({success:false,message:'需要超级管理员会话'});await current.fetch({useMasterKey:true});
- 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:'仅超级管理员可执行'});
- const companyId=${JSON.stringify(COMPANY_ID)},merge=${embeddedMerge},promote=${embeddedPromote};
- await Psql.none('ALTER TABLE "SurveyLog" ALTER COLUMN "user" TYPE text USING CAST("user" AS text)');
- 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};
- 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};
- response.json({success:true,data:{merged,promoted}});
- }catch(error){response.status(Number(error.status)||500).json({success:false,message:String(error.message||error)});}
- }`;
- const fn = await parse('/classes/Function', {
- method: 'POST',
- body: JSON.stringify({
- name: path, desc: '一次性合并阅读账号与原会员', type: 'standalone', path, code,
- params: [], paramList: [], respType: 'json', respJson: { success: true }, enabled: true,
- }),
- });
- temporaryFunctionId = fn.objectId;
- const response = await fetch(`${FUNCTION_URL}/${path}`, {
- method: 'POST',
- headers: { 'X-Parse-Application-Id': APP_ID, 'Content-Type': 'application/json' },
- body: JSON.stringify({ token, params: {} }),
- });
- const payload = await response.json().catch(() => ({}));
- if (!response.ok || payload.success !== true) throw new Error(payload.message || payload.error || `账号合并失败:${response.status}`);
- return payload.data;
- }
- async function cleanup() {
- if (temporaryFunctionId) await parse(`/classes/Function/${temporaryFunctionId}`, { method: 'DELETE' }).catch(() => undefined);
- if (temporaryUserId) {
- const where = encodeURIComponent(JSON.stringify({ user: { __type: 'Pointer', className: '_User', objectId: temporaryUserId } }));
- const sessions = await parse(`/classes/_Session?where=${where}&limit=1000&keys=objectId`).catch(() => ({ results: [] }));
- for (const session of sessions.results || []) await parse(`/classes/_Session/${session.objectId}`, { method: 'DELETE' }).catch(() => undefined);
- await parse(`/users/${temporaryUserId}`, { method: 'DELETE' }).catch(() => undefined);
- }
- }
- const userKeys = 'objectId,username,nickname,nickName,realName,name,mobile,phone,mobilePhoneNumber,level,company,sourceKey,legacyUserId,legacyGroupId,legacyUserData,identityType,isAdmin,isDisabled,isDeleted,createdAt,registeredAt';
- const usersBefore = await allRows('/users', { company }, userKeys);
- const ownersBefore = await ownerStats();
- const plan = buildPlan(usersBefore);
- const sourceIds = new Set(plan.sources.map((row) => row.objectId));
- const readingSources = new Set();
- for (const value of Object.values(ownersBefore)) {
- for (const owner of value.counts.keys()) if (sourceIds.has(owner)) readingSources.add(owner);
- }
- const summary = {
- mode: commit ? 'commit' : 'dry-run',
- usersBefore: usersBefore.length,
- canonicalUsersBefore: plan.targets.length,
- readingImports: plan.sources.length,
- readingImportsWithRecords: readingSources.size,
- mergeUsers: plan.merge.length,
- promoteUsers: plan.promote.length,
- maxMergeSourceIdLength: Math.max(0, ...plan.merge.map((item) => item.source.length)),
- maxMergeTargetIdLength: Math.max(0, ...plan.merge.map((item) => item.target.length)),
- nameResolved: plan.nameResolved,
- ambiguousPromoted: plan.ambiguous,
- unmatchedPromoted: plan.unmatched,
- surveyItemsToMove: plan.merge.reduce((total, item) => total + (ownersBefore.SurveyItem.counts.get(item.source) || 0), 0),
- surveyLogsToMove: plan.merge.reduce((total, item) => total + (ownersBefore.SurveyLog.counts.get(item.source) || 0), 0),
- };
- try {
- if (commit && (plan.merge.length || plan.promote.length)) summary.executed = await executePlan(plan);
- if (commit) {
- const usersAfter = await allRows('/users', { company }, userKeys);
- const ownersAfter = await ownerStats();
- const remaining = usersAfter.filter(isReadingImport);
- const activeById = new Map(usersAfter.filter((row) => !row.isDeleted).map((row) => [row.objectId, row]));
- const invalidOwners = new Set();
- for (const value of Object.values(ownersAfter)) {
- for (const owner of value.counts.keys()) {
- const row = activeById.get(owner);
- if (!row || (!isCanonical(row) && !text(row.sourceKey).startsWith('reading-unified:user:'))) invalidOwners.add(owner);
- }
- }
- summary.verification = {
- usersAfter: usersAfter.length,
- remainingUnnormalizedUsers: remaining.length,
- invalidReadingOwners: invalidOwners.size,
- surveyItems: ownersAfter.SurveyItem.records,
- surveyLogs: ownersAfter.SurveyLog.records,
- };
- if (remaining.length || invalidOwners.size) throw new Error(`合并后验证失败:未归一账号 ${remaining.length},无效阅读数据归属 ${invalidOwners.size}`);
- }
- process.stdout.write(`${JSON.stringify(summary, null, 2)}\n`);
- } finally {
- await cleanup();
- }
|