Files
koc-loop/app/api/resources-import/route.ts
2026-08-15 03:53:09 +08:00

501 lines
17 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 });
}
}