import { isManagerRequest, managerForbidden } from "../../../lib/user-auth"; import { runInBackground } from "../../../lib/background"; import { ensureSchema, getRawDb } from "../../../lib/mvp-db"; import { resolveCollectionMcpConfig, resolveProfileDetailsFromMcp, resolveXhsPublicAccountDetails, type CollectionMcpBindings, } from "../../../lib/mcp-collection-client"; import { getRuntimeEnv } from "../../../lib/runtime-env"; import { mergeCooperationSources, normalizeProfileUrl, parseResourceImportFile, RESOURCE_IMPORT_MAX_BYTES, RESOURCE_IMPORT_MAX_ROWS, resourcePlatformUid, resourceImportMissingFields, type ResourceImportRow, } from "../../../lib/resource-import"; type AccountRow = { id: string; platform: string; platform_uid: string; public_account_id: string; nickname: string; profile_url: string; ip_location: string; followers: number; gender: string; bio: string; tags: string; cooperation_source: string; }; type AnalyzedRow = ResourceImportRow & { action: "create" | "update" | "error"; accountId: string; platformUid: string; cooperationSource: string; }; const RESOURCE_IMPORT_PREVIEW_ROWS = 100; const RESOURCE_IMPORT_SYNC_ENRICH_ROWS = 100; const RESOURCE_IMPORT_DB_BATCH_SIZE = 100; function identityKey(platform: string, value: string) { return `${platform.trim().toLocaleLowerCase("zh-CN")}|${value .trim() .toLocaleLowerCase("zh-CN")}`; } async function loadAccounts() { return getRawDb() .prepare( `SELECT id, platform, platform_uid, public_account_id, nickname, profile_url, ip_location, followers, gender, bio, tags, cooperation_source FROM accounts`, ) .all(); } async function mapConcurrent( items: T[], limit: number, worker: (item: T) => Promise, ) { const results = new Array(items.length); let cursor = 0; await Promise.all( Array.from({ length: Math.min(limit, items.length) }, async () => { while (cursor < items.length) { const index = cursor; cursor += 1; results[index] = await worker(items[index]); } }), ); return results; } function mergeExistingFields(rows: ResourceImportRow[], accounts: AccountRow[]) { const existingByProfile = new Map(); for (const account of accounts) { const profileUrl = normalizeProfileUrl(account.profile_url || ""); if (profileUrl) { existingByProfile.set(identityKey(account.platform, profileUrl), account); } } return rows.map((row) => { if ( row.errors.length > 0 || !row.profileUrl || !["小红书", "抖音"].includes(row.platform) ) { return row; } const existing = existingByProfile.get(identityKey(row.platform, row.profileUrl)); const existingNickname = existing?.nickname && existing.nickname !== "待识别账号" ? existing.nickname : ""; const existingIpLocation = existing?.ip_location && existing.ip_location !== "待识别" ? existing.ip_location : ""; const existingGender: ResourceImportRow["gender"] = existing?.gender === "男" || existing?.gender === "女" ? existing.gender : ""; return { ...row, nickname: row.nickname || existingNickname, publicAccountId: row.publicAccountId || existing?.public_account_id || "", ipLocation: row.ipLocation || existingIpLocation, followers: row.followersResolved ? row.followers : Number(existing?.followers || 0), followersResolved: row.followersResolved || Number(existing?.followers || 0) > 0, gender: row.gender || existingGender, bio: row.bio || existing?.bio || "", tags: row.tags.length > 0 ? row.tags : (existing?.tags || "") .split(/[,,、;;|]/) .map((item) => item.trim()) .filter(Boolean) .slice(0, 5), }; }); } async function enrichRows(rows: ResourceImportRow[], accounts: AccountRow[]) { const mcpConfig = resolveCollectionMcpConfig( getRuntimeEnv() as unknown as CollectionMcpBindings, ); const baselineRows = mergeExistingFields(rows, accounts); return mapConcurrent(baselineRows, 4, async (baseline) => { const row = baseline; if ( row.errors.length > 0 || !row.profileUrl || !["小红书", "抖音"].includes(row.platform) ) { return row; } if (resourceImportMissingFields(baseline).length === 0) { return baseline; } let details: { nickname: string | null; redId: string | null; followers: number | null; ipLocation: string | null; gender: "" | "男" | "女"; bio: string; recentNoteTitles: string[]; providerTags: string[]; } = await resolveProfileDetailsFromMcp( row.profileUrl, row.platform === "抖音" ? "抖音" : "小红书", mcpConfig, ).catch(() => ({ nickname: null, redId: null, followers: null, ipLocation: null, gender: "" as const, bio: "", recentNoteTitles: [], providerTags: [], })); const mcpResult = { nickname: baseline.nickname || details.nickname?.trim() || "", publicAccountId: baseline.publicAccountId || details.redId?.trim() || "", ipLocation: baseline.ipLocation || details.ipLocation?.trim() || "", followersResolved: baseline.followersResolved || details.followers !== null, gender: baseline.gender || details.gender, bio: baseline.bio || details.bio, tags: baseline.tags, }; if ( row.platform === "小红书" && resourceImportMissingFields(mcpResult).length > 0 ) { const publicDetails = await resolveXhsPublicAccountDetails(row.profileUrl).catch( () => ({ nickname: null, redId: null, followers: null, ipLocation: null }), ); details = { ...details, nickname: details.nickname || publicDetails.nickname, redId: details.redId || publicDetails.redId, followers: details.followers ?? publicDetails.followers, ipLocation: details.ipLocation || publicDetails.ipLocation, }; } return { ...baseline, nickname: baseline.nickname || details.nickname?.trim() || "待识别账号", publicAccountId: baseline.publicAccountId || details.redId?.trim() || "", ipLocation: baseline.ipLocation || details.ipLocation?.trim() || "待识别", followers: baseline.followersResolved ? baseline.followers : (details.followers ?? 0), followersResolved: baseline.followersResolved || details.followers !== null, gender: baseline.gender || details.gender, bio: baseline.bio || details.bio, tags: baseline.tags, }; }); } function analyzeRows(rows: ResourceImportRow[], accountRows: AccountRow[]) { const profileMap = new Map(); const publicIdMap = new Map(); const platformUidMap = new Map(); for (const account of accountRows) { const profileUrl = normalizeProfileUrl(account.profile_url || ""); if (profileUrl) profileMap.set(identityKey(account.platform, profileUrl), account); if (account.public_account_id) { publicIdMap.set(identityKey(account.platform, account.public_account_id), account); } platformUidMap.set(identityKey(account.platform, account.platform_uid), account); } return rows.map((row) => { const platformUid = resourcePlatformUid(row); const profileMatch = row.profileUrl ? profileMap.get(identityKey(row.platform, row.profileUrl)) : undefined; const publicIdMatch = row.publicAccountId ? publicIdMap.get(identityKey(row.platform, row.publicAccountId)) : undefined; const uidMatch = platformUidMap.get(identityKey(row.platform, platformUid)); const matches = [profileMatch, publicIdMatch, uidMatch].filter( (account): account is AccountRow => Boolean(account), ); const matchedIds = [...new Set(matches.map((account) => account.id))]; const errors = [...row.errors]; if (matchedIds.length > 1) { errors.push("账号主页和账号ID匹配到不同的现有账号,请先核对"); } const existing = matchedIds.length === 1 ? matches[0] : undefined; const accountId = existing?.id ?? `account-${crypto.randomUUID().slice(0, 12)}`; const analyzed: AnalyzedRow = { ...row, errors, action: errors.length > 0 ? "error" : existing ? "update" : "create", accountId, platformUid: existing?.platform_uid ?? platformUid, cooperationSource: mergeCooperationSources( existing?.cooperation_source ?? "", row.cooperationSource, ), }; if (analyzed.action !== "error") { const virtual: AccountRow = { id: accountId, platform: row.platform, platform_uid: analyzed.platformUid, public_account_id: row.publicAccountId || existing?.public_account_id || "", nickname: row.nickname, profile_url: row.profileUrl || existing?.profile_url || "", ip_location: row.ipLocation || existing?.ip_location || "待识别", followers: row.followers || existing?.followers || 0, gender: row.gender || existing?.gender || "", bio: row.bio || existing?.bio || "", tags: (row.tags.length > 0 ? row.tags : (existing?.tags || "").split(/[,,、;;|]/).filter(Boolean) ).slice(0, 5).join(","), cooperation_source: analyzed.cooperationSource, }; if (row.profileUrl) profileMap.set(identityKey(row.platform, row.profileUrl), virtual); if (row.publicAccountId) { publicIdMap.set(identityKey(row.platform, row.publicAccountId), virtual); } platformUidMap.set(identityKey(row.platform, analyzed.platformUid), virtual); } return analyzed; }); } function summarize(rows: AnalyzedRow[]) { return { total: rows.length, create: rows.filter((row) => row.action === "create").length, update: rows.filter((row) => row.action === "update").length, error: rows.filter((row) => row.action === "error").length, }; } function previewAnalyzedRows(rows: AnalyzedRow[]) { const errorRows = rows.filter((row) => row.action === "error"); if (errorRows.length === 0) { return rows.slice(0, RESOURCE_IMPORT_PREVIEW_ROWS); } const importableRows = rows.filter((row) => row.action !== "error"); return [ ...errorRows.slice(0, RESOURCE_IMPORT_PREVIEW_ROWS), ...importableRows.slice( 0, Math.max(0, RESOURCE_IMPORT_PREVIEW_ROWS - errorRows.length), ), ]; } function statementForAnalyzedRow( db: ReturnType, row: AnalyzedRow, ) { return row.action === "update" ? db .prepare( `UPDATE accounts SET nickname = ?, public_account_id = CASE WHEN ? != '' THEN ? ELSE public_account_id END, profile_url = CASE WHEN ? != '' THEN ? ELSE profile_url END, ip_location = CASE WHEN ? != '' AND ? != '待识别' THEN ? ELSE ip_location END, followers = CASE WHEN ? = 1 THEN ? ELSE followers END, gender = CASE WHEN ? != '' THEN ? ELSE gender END, bio = CASE WHEN ? != '' THEN ? ELSE bio END, tags = CASE WHEN ? != '' THEN ? ELSE tags END, cooperation_source = ?, last_seen_at = CURRENT_TIMESTAMP WHERE id = ?`, ) .bind( row.nickname || "待识别账号", row.publicAccountId, row.publicAccountId, row.profileUrl, row.profileUrl, row.ipLocation, row.ipLocation, row.ipLocation, row.followersResolved ? 1 : 0, row.followers, row.gender, row.gender, row.bio, row.bio, row.tags.join(","), row.tags.join(","), row.cooperationSource, row.accountId, ) : db .prepare( `INSERT INTO accounts (id, platform, platform_uid, public_account_id, nickname, profile_url, ip_location, followers, post_count, avg_views, gender, bio, tags, cooperation_source) VALUES (?, ?, ?, ?, ?, ?, ?, ?, 0, 0, ?, ?, ?, ?)`, ) .bind( row.accountId, row.platform, row.platformUid, row.publicAccountId, row.nickname || "待识别账号", row.profileUrl, row.ipLocation || "待识别", row.followers, row.gender, row.bio, row.tags.join(","), row.cooperationSource, ); } async function writeAnalyzedRows(rows: AnalyzedRow[]) { const db = getRawDb(); const statements = rows .filter((row) => row.action !== "error") .map((row) => statementForAnalyzedRow(db, row)); for (let index = 0; index < statements.length; index += RESOURCE_IMPORT_DB_BATCH_SIZE) { await db.batch(statements.slice(index, index + RESOURCE_IMPORT_DB_BATCH_SIZE)); } } function deferredEnrichmentCount(rows: ResourceImportRow[]) { return rows.filter( (row) => row.errors.length === 0 && row.profileUrl && resourceImportMissingFields(row).length > 0, ).length; } async function enrichImportedRowsInBackground(rows: ResourceImportRow[]) { const accounts = await loadAccounts(); const baselineRows = mergeExistingFields(rows, accounts.results); const missingRows = baselineRows.filter( (row) => row.errors.length === 0 && row.profileUrl && resourceImportMissingFields(row).length > 0, ); if (missingRows.length === 0) return; const enriched = await enrichRows(missingRows, accounts.results); const latestAccounts = await loadAccounts(); const analyzed = analyzeRows(enriched, latestAccounts.results).filter( (row) => row.action !== "error", ); await writeAnalyzedRows(analyzed); } export async function POST(request: Request) { if (!(await isManagerRequest(request))) return managerForbidden(); try { await ensureSchema(); const form = await request.formData(); const file = form.get("file"); const mode = String(form.get("mode") ?? "preview"); if (!(file instanceof File)) { return Response.json({ error: "请选择需要导入的 Excel 文件" }, { status: 400 }); } if (file.size <= 0 || file.size > RESOURCE_IMPORT_MAX_BYTES) { return Response.json({ error: "文件不能为空,且不能超过 20MB" }, { status: 400 }); } const rows = parseResourceImportFile(file.name, new Uint8Array(await file.arrayBuffer())); const accounts = await loadAccounts(); const shouldEnrichSynchronously = rows.length <= RESOURCE_IMPORT_SYNC_ENRICH_ROWS; const preparedRows = shouldEnrichSynchronously ? await enrichRows(rows, accounts.results) : mergeExistingFields(rows, accounts.results); const analyzed = analyzeRows(preparedRows, accounts.results); const summary = summarize(analyzed); const importableRows = analyzed.filter((row) => row.action !== "error"); const importableRowNumbers = new Set( importableRows.map((row) => row.rowNumber), ); const importablePreparedRows = preparedRows.filter((row) => importableRowNumbers.has(row.rowNumber), ); const deferredEnrichment = shouldEnrichSynchronously ? 0 : deferredEnrichmentCount(importablePreparedRows); if (mode !== "commit") { return Response.json({ summary, rows: previewAnalyzedRows(analyzed).map((row) => ({ rowNumber: row.rowNumber, platform: row.platform, nickname: row.nickname, publicAccountId: row.publicAccountId, profileUrl: row.profileUrl, ipLocation: row.ipLocation, followers: row.followers, gender: row.gender, bio: row.bio, tags: row.tags, cooperationSource: row.cooperationSource, action: row.action, errors: row.errors, })), truncated: analyzed.length > RESOURCE_IMPORT_PREVIEW_ROWS, deferredEnrichment, maxRows: RESOURCE_IMPORT_MAX_ROWS, }); } if (importableRows.length === 0) { return Response.json( { error: "没有可导入的有效账号,请修正异常数据后重新上传", summary }, { status: 400 }, ); } await writeAnalyzedRows(importableRows); if (deferredEnrichment > 0) { const importableSourceRows = rows.filter((row) => importableRowNumbers.has(row.rowNumber), ); runInBackground( enrichImportedRowsInBackground(importableSourceRows), "bulk resource profile enrichment", ); } return Response.json({ summary, deferredEnrichment, message: `已导入 ${summary.create} 个新账号,更新 ${summary.update} 个已有账号${ summary.error > 0 ? `;跳过 ${summary.error} 条异常数据` : "" }${ deferredEnrichment > 0 ? `;${deferredEnrichment} 个账号的缺失公开资料将在后台补全` : "" }`, }); } catch (error) { const message = error instanceof Error ? error.message : "导入失败"; return Response.json({ error: message }, { status: 400 }); } }