Files
koc-loop/lib/mcp-operations.ts
2026-08-12 11:43:26 +08:00

465 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 { enrichDistributionAccount } from "./account-enrichment-service";
import {
collectDistributionMetrics,
createCollectionRunTasks,
retryFailedCollections,
runDueScheduledCollections,
shanghaiDateFromTimestamp,
} from "./collection-service";
import { issueMcpExportToken } from "./mcp-export-token";
import {
resolveCollectionMcpConfig,
type CollectionMcpBindings,
} from "./mcp-collection-client";
import { ensureSchema, getRawDb } from "./mvp-db";
import { extractXhsPublishUrl } from "./publish-url";
import { buildClaimUrl } from "./task-service";
export type McpOperationBindings = CollectionMcpBindings & {
KOC_PORTAL_URL?: string;
};
type Pagination = { limit?: number; offset?: number };
type ResourceFilters = Pagination & {
query?: string;
ipLocation?: string;
cooperationSource?: string;
platform?: string;
};
function page(input: Pagination) {
return {
limit: Math.max(1, Math.min(200, Math.floor(input.limit ?? 50))),
offset: Math.max(0, Math.floor(input.offset ?? 0)),
};
}
function like(value: string) {
return `%${value.replace(/[\\%_]/g, "\\$&")}%`;
}
function parseJsonArray(value: unknown) {
try {
const parsed = JSON.parse(String(value ?? "[]"));
return Array.isArray(parsed) ? parsed : [];
} catch {
return [];
}
}
async function taskExists(taskId: string) {
const task = await getRawDb()
.prepare("SELECT id FROM tasks WHERE id = ?")
.bind(taskId)
.first<{ id: string }>();
if (!task) throw new Error("任务不存在");
}
export async function taskList(
input: Pagination & { query?: string; status?: string },
portalUrl: string,
) {
await ensureSchema();
const { limit, offset } = page(input);
const conditions: string[] = [];
const bindings: unknown[] = [];
if (input.query?.trim()) {
conditions.push("(t.name LIKE ? ESCAPE '\\' OR t.brand LIKE ? ESCAPE '\\')");
const pattern = like(input.query.trim());
bindings.push(pattern, pattern);
}
if (input.status?.trim() && input.status !== "all") {
conditions.push("t.status = ?");
bindings.push(input.status.trim());
}
const where = conditions.length ? `WHERE ${conditions.join(" AND ")}` : "";
const db = getRawDb();
const [rows, count] = await Promise.all([
db
.prepare(
`SELECT t.*,
(SELECT COUNT(*) FROM contents c WHERE c.task_id = t.id) AS note_count,
(SELECT COUNT(*) FROM distributions d WHERE d.task_id = t.id) AS claimed_count,
(SELECT COUNT(*) FROM distributions d WHERE d.task_id = t.id AND COALESCE(d.publish_url, '') != '') AS published_count,
(SELECT COUNT(*) FROM distributions d WHERE d.task_id = t.id AND d.exposure IS NOT NULL AND d.views IS NOT NULL) AS day7_count
FROM tasks t ${where}
ORDER BY t.created_at DESC LIMIT ? OFFSET ?`,
)
.bind(...bindings, limit, offset)
.all<Record<string, unknown>>(),
db
.prepare(`SELECT COUNT(*) AS total FROM tasks t ${where}`)
.bind(...bindings)
.first<{ total: number }>(),
]);
return {
total: count?.total ?? 0,
limit,
offset,
tasks: rows.results.map((row) => ({
...row,
collection_days: parseJsonArray(row.collection_days),
claim_url:
portalUrl && row.share_token
? buildClaimUrl(portalUrl, String(row.share_token))
: null,
})),
};
}
export async function taskGet(taskId: string, portalUrl: string) {
await ensureSchema();
const db = getRawDb();
const task = await db
.prepare("SELECT * FROM tasks WHERE id = ?")
.bind(taskId)
.first<Record<string, unknown>>();
if (!task) throw new Error("任务不存在");
const [notes, claims, runs] = await Promise.all([
db
.prepare(
`SELECT c.id AS content_id, c.source_row, c.title, c.body, c.image_assets, c.status AS content_status,
d.id AS distribution_id, d.status AS distribution_status, d.publish_url, d.publish_time,
d.exposure, d.views, d.latest_likes, d.latest_comments, d.latest_collects,
d.collection_status, d.collection_status_description, d.collection_updated_at,
a.id AS account_id, a.nickname AS account_nickname, a.profile_url, a.followers,
p.name AS cooperation_source, cl.claimant_name
FROM contents c
LEFT JOIN distributions d ON d.content_id = c.id
LEFT JOIN accounts a ON a.id = d.account_id
LEFT JOIN partners p ON p.id = d.partner_id
LEFT JOIN claims cl ON cl.id = d.claim_id
WHERE c.task_id = ?
ORDER BY COALESCE(c.source_row, 999999), c.created_at, d.claimed_at`,
)
.bind(taskId)
.all<Record<string, unknown>>(),
db
.prepare(
`SELECT cl.*, p.name AS partner_name,
(SELECT COUNT(*) FROM distributions d WHERE d.claim_id = cl.id AND COALESCE(d.publish_url, '') != '') AS published_count
FROM claims cl JOIN partners p ON p.id = cl.partner_id
WHERE cl.task_id = ? ORDER BY cl.created_at DESC`,
)
.bind(taskId)
.all<Record<string, unknown>>(),
db
.prepare(
`SELECT * FROM collection_runs WHERE task_id = ?
ORDER BY scheduled_date DESC, created_at DESC LIMIT 500`,
)
.bind(taskId)
.all<Record<string, unknown>>(),
]);
return {
task: {
...task,
collection_days: parseJsonArray(task.collection_days),
claim_url:
portalUrl && task.share_token
? buildClaimUrl(portalUrl, String(task.share_token))
: null,
},
notes: notes.results.map((row) => ({
...row,
image_assets: parseJsonArray(row.image_assets),
total_interactions:
row.latest_likes == null
? null
: Number(row.latest_likes) +
Number(row.latest_comments ?? 0) +
Number(row.latest_collects ?? 0),
})),
claims: claims.results,
collection_runs: runs.results,
};
}
export async function recoveryList(
input: Pagination & {
taskId?: string;
stage?: "all" | "published" | "unfilled" | "waiting_day7" | "day7_due";
},
) {
await ensureSchema();
const { limit, offset } = page(input);
const db = getRawDb();
const conditions: string[] = [];
const bindings: unknown[] = [];
if (input.taskId?.trim()) {
conditions.push("d.task_id = ?");
bindings.push(input.taskId.trim());
}
const stage = input.stage ?? "all";
if (stage === "published") conditions.push("COALESCE(d.publish_url, '') != ''");
if (stage === "unfilled") conditions.push("COALESCE(d.publish_url, '') = ''");
if (stage === "waiting_day7") {
conditions.push("COALESCE(d.publish_url, '') != '' AND (d.exposure IS NULL OR d.views IS NULL)");
}
if (stage === "day7_due") {
conditions.push(
"COALESCE(d.publish_url, '') != '' AND (d.exposure IS NULL OR d.views IS NULL) AND datetime(d.publish_time, '+7 days') <= CURRENT_TIMESTAMP",
);
}
const where = conditions.length ? `WHERE ${conditions.join(" AND ")}` : "";
const base = `FROM distributions d
JOIN tasks t ON t.id = d.task_id
JOIN contents c ON c.id = d.content_id
LEFT JOIN accounts a ON a.id = d.account_id
LEFT JOIN partners p ON p.id = d.partner_id
LEFT JOIN claims cl ON cl.id = d.claim_id ${where}`;
const [rows, count] = await Promise.all([
db
.prepare(
`SELECT d.*, t.name AS task_name, c.source_row, c.title,
a.nickname AS account_nickname, a.profile_url, p.name AS cooperation_source,
cl.claimant_name ${base}
ORDER BY COALESCE(d.publish_time, d.claimed_at) DESC LIMIT ? OFFSET ?`,
)
.bind(...bindings, limit, offset)
.all<Record<string, unknown>>(),
db
.prepare(`SELECT COUNT(*) AS total ${base}`)
.bind(...bindings)
.first<{ total: number }>(),
]);
return {
total: count?.total ?? 0,
limit,
offset,
items: rows.results.map((row) => ({
...row,
total_interactions:
row.latest_likes == null
? null
: Number(row.latest_likes) + Number(row.latest_comments ?? 0) + Number(row.latest_collects ?? 0),
})),
};
}
export async function recoveryExport(taskId: string, origin: string) {
await ensureSchema();
await taskExists(taskId);
const issued = await issueMcpExportToken("recovery", { taskId });
return {
task_id: taskId,
download_url: `${origin}/api/recovery-export?token=${encodeURIComponent(issued.token)}`,
expires_at: issued.expiresAt,
};
}
export async function setCollectionPlan(
taskId: string,
startDate: string,
days: number[],
bindings: McpOperationBindings,
) {
await ensureSchema();
await taskExists(taskId);
const normalizedDays = [...new Set(days)]
.filter((day) => Number.isInteger(day) && day >= 1 && day <= 7)
.sort((a, b) => a - b);
if (!/^\d{4}-\d{2}-\d{2}$/.test(startDate) || Number.isNaN(Date.parse(`${startDate}T00:00:00+08:00`))) {
throw new Error("开始采集日期无效");
}
if (!normalizedDays.length) throw new Error("请至少选择一个采集日");
const db = getRawDb();
await db.batch([
db.prepare(
`UPDATE tasks SET collection_start_date = ?, collection_days = ?,
collection_schedule_updated_at = CURRENT_TIMESTAMP WHERE id = ?`,
).bind(startDate, JSON.stringify(normalizedDays), taskId),
db.prepare(
`UPDATE distributions SET
collection_status = CASE WHEN latest_likes IS NULL THEN 'scheduled' ELSE collection_status END,
collection_status_description = CASE WHEN latest_likes IS NULL THEN ? ELSE collection_status_description END,
updated_at = CURRENT_TIMESTAMP
WHERE task_id = ? AND COALESCE(publish_url, '') != ''`,
).bind(`已安排${normalizedDays.length}个采集日每日09:00执行`, taskId),
]);
const created = await createCollectionRunTasks(db, taskId, startDate, normalizedDays);
const catchup = await runDueScheduledCollections(
db,
Date.now(),
resolveCollectionMcpConfig(bindings),
"catchup",
taskId,
);
return { task_id: taskId, start_date: startDate, days: normalizedDays, created, catchup };
}
export async function runDueCollections(taskId: string | undefined, bindings: McpOperationBindings) {
await ensureSchema();
if (taskId) await taskExists(taskId);
return runDueScheduledCollections(
getRawDb(),
Date.now(),
resolveCollectionMcpConfig(bindings),
"catchup",
taskId,
);
}
export async function collectNow(
distributionId: string,
scheduleDay: number | undefined,
bindings: McpOperationBindings,
) {
await ensureSchema();
const day = scheduleDay ?? null;
const metrics = await collectDistributionMetrics(
getRawDb(),
distributionId,
shanghaiDateFromTimestamp(Date.now()),
day,
"manual",
resolveCollectionMcpConfig(bindings),
);
return { distribution_id: distributionId, schedule_day: day, ...metrics };
}
export async function retryFailed(taskId: string | undefined, bindings: McpOperationBindings) {
await ensureSchema();
const db = getRawDb();
const ids = taskId
? [taskId]
: (
await db
.prepare(
`SELECT DISTINCT task_id FROM distributions
WHERE collection_status = 'failed' ORDER BY updated_at DESC LIMIT 20`,
)
.all<{ task_id: string }>()
).results.map((row) => row.task_id);
if (taskId) await taskExists(taskId);
const results = [];
for (const id of ids) {
results.push({ task_id: id, ...(await retryFailedCollections(db, id, resolveCollectionMcpConfig(bindings))) });
}
return {
task_count: results.length,
attempted: results.reduce((sum, item) => sum + item.attempted, 0),
succeeded: results.reduce((sum, item) => sum + item.succeeded, 0),
failed: results.reduce((sum, item) => sum + item.failed, 0),
tasks: results,
};
}
function resourceWhere(input: ResourceFilters) {
const conditions: string[] = [];
const bindings: unknown[] = [];
if (input.query?.trim()) {
conditions.push("(a.nickname LIKE ? ESCAPE '\\' OR a.public_account_id LIKE ? ESCAPE '\\')");
const pattern = like(input.query.trim());
bindings.push(pattern, pattern);
}
if (input.ipLocation?.trim()) {
conditions.push("a.ip_location LIKE ? ESCAPE '\\'");
bindings.push(like(input.ipLocation.trim()));
}
if (input.cooperationSource?.trim()) {
conditions.push(
`EXISTS (SELECT 1 FROM distributions dx JOIN partners px ON px.id = dx.partner_id
WHERE dx.account_id = a.id AND px.name LIKE ? ESCAPE '\\')`,
);
bindings.push(like(input.cooperationSource.trim()));
}
if (input.platform?.trim() && input.platform !== "all") {
conditions.push("a.platform = ?");
bindings.push(input.platform.trim());
}
return { where: conditions.length ? `WHERE ${conditions.join(" AND ")}` : "", bindings };
}
export async function resourceSearch(input: ResourceFilters) {
await ensureSchema();
const { limit, offset } = page(input);
const { where, bindings } = resourceWhere(input);
const db = getRawDb();
const select = `SELECT a.*,
(SELECT COUNT(*) FROM distributions d WHERE d.account_id = a.id) AS cooperation_count,
(SELECT GROUP_CONCAT(DISTINCT p.name) FROM distributions d JOIN partners p ON p.id = d.partner_id WHERE d.account_id = a.id) AS cooperation_sources`;
const [rows, count] = await Promise.all([
db.prepare(`${select} FROM accounts a ${where} ORDER BY a.last_seen_at DESC LIMIT ? OFFSET ?`)
.bind(...bindings, limit, offset).all<Record<string, unknown>>(),
db.prepare(`SELECT COUNT(*) AS total FROM accounts a ${where}`).bind(...bindings).first<{ total: number }>(),
]);
return {
total: count?.total ?? 0,
limit,
offset,
accounts: rows.results.map((row) => ({
...row,
cooperation_sources: String(row.cooperation_sources ?? "").split(",").filter(Boolean),
})),
};
}
export async function resourceGet(accountId: string) {
await ensureSchema();
const db = getRawDb();
const account = await db.prepare("SELECT * FROM accounts WHERE id = ?").bind(accountId).first<Record<string, unknown>>();
if (!account) throw new Error("账号不存在");
const history = await db.prepare(
`SELECT d.id AS distribution_id, d.task_id, t.name AS task_name, c.title,
d.publish_url, d.publish_time, d.latest_likes, d.latest_comments, d.latest_collects,
d.exposure, d.views, p.name AS cooperation_source, cl.claimant_name
FROM distributions d JOIN tasks t ON t.id = d.task_id
JOIN contents c ON c.id = d.content_id JOIN partners p ON p.id = d.partner_id
LEFT JOIN claims cl ON cl.id = d.claim_id
WHERE d.account_id = ? ORDER BY COALESCE(d.publish_time, d.claimed_at) DESC`,
).bind(accountId).all<Record<string, unknown>>();
return { account, cooperation_history: history.results };
}
export async function backfillResourceProfile(
input: { distributionId?: string; publishUrl?: string },
bindings: McpOperationBindings,
) {
await ensureSchema();
const db = getRawDb();
const publishUrl = input.publishUrl ? extractXhsPublishUrl(input.publishUrl) : "";
const row = input.distributionId
? await db.prepare(
`SELECT d.id, d.publish_url, COALESCE(a.nickname, '') AS nickname
FROM distributions d LEFT JOIN accounts a ON a.id = d.account_id WHERE d.id = ?`,
).bind(input.distributionId).first<{ id: string; publish_url: string | null; nickname: string }>()
: publishUrl
? await db.prepare(
`SELECT d.id, d.publish_url, COALESCE(a.nickname, '') AS nickname
FROM distributions d LEFT JOIN accounts a ON a.id = d.account_id
WHERE d.publish_url = ? ORDER BY d.updated_at DESC LIMIT 1`,
).bind(publishUrl).first<{ id: string; publish_url: string | null; nickname: string }>()
: null;
if (!row?.publish_url) throw new Error("没有找到可补全的发布记录");
const result = await enrichDistributionAccount(
db,
row.id,
row.publish_url,
row.nickname,
resolveCollectionMcpConfig(bindings),
);
const account = result.updated && result.accountId
? await db.prepare("SELECT * FROM accounts WHERE id = ?").bind(result.accountId).first<Record<string, unknown>>()
: null;
return { distribution_id: row.id, ...result, account };
}
export async function resourceExport(input: ResourceFilters, origin: string) {
const result = await resourceSearch({ ...input, limit: 200, offset: 0 });
const allIds: string[] = result.accounts.map((account) => String(account.id));
if (result.total > result.accounts.length) {
const { where, bindings } = resourceWhere(input);
const rows = await getRawDb().prepare(`SELECT a.id FROM accounts a ${where} ORDER BY a.last_seen_at DESC LIMIT 5000`)
.bind(...bindings).all<{ id: string }>();
allIds.splice(0, allIds.length, ...rows.results.map((row) => row.id));
}
if (!allIds.length) throw new Error("当前筛选结果为空");
const issued = await issueMcpExportToken("resources", { accountIds: allIds });
return {
account_count: allIds.length,
download_url: `${origin}/api/resources-export?token=${encodeURIComponent(issued.token)}`,
expires_at: issued.expiresAt,
};
}