501 lines
17 KiB
TypeScript
501 lines
17 KiB
TypeScript
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<AccountRow>();
|
||
}
|
||
|
||
async function mapConcurrent<T, R>(
|
||
items: T[],
|
||
limit: number,
|
||
worker: (item: T) => Promise<R>,
|
||
) {
|
||
const results = new Array<R>(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<string, AccountRow>();
|
||
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<string, AccountRow>();
|
||
const publicIdMap = new Map<string, AccountRow>();
|
||
const platformUidMap = new Map<string, AccountRow>();
|
||
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<AnalyzedRow>((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<typeof getRawDb>,
|
||
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 });
|
||
}
|
||
}
|