#!/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(/>\'{}\' 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(); }