|
|
@@ -0,0 +1,256 @@
|
|
|
+#!/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();
|
|
|
+}
|