3 Commits

Author SHA1 Message Date
巫凤萍
ab3a544eec feat: add MCP creator screenshot metrics recognition 2026-08-19 15:08:17 +08:00
8f7ea0558d Merge pull request 'feat: 接入企业微信通知(临期催办 + 群机器人汇总)' (#3) from feat/account-tags into main
Reviewed-on: #3
2026-08-18 09:12:59 +00:00
ABAPPLO
f2ac751c4c feat: 接入企业微信通知(临期催办 + 群机器人汇总)
新增 lib/wecom-client.ts 与 lib/wecom-notifier-service.ts,支持
群机器人 webhook 与应用消息双通道;scheduler 加入临期 N 天催办
与每日管理员汇总;action 路由补 bind_wecom_external_id 与
send_test_wecom 两个管理端动作,配套测试。
2026-08-18 17:05:16 +08:00
20 changed files with 1071 additions and 38 deletions

View File

@@ -24,6 +24,15 @@ FEISHU_APP_SECRET=
AI_TOOL_CENTER_MCP_URL=
AI_TOOL_CENTER_MCP_KEY=
# 企业微信通知(临期催办 + 管理员汇总)。群机器人只需 webhookKOC 侧催办还需 corp/agent/secret
# 并在后台 partners 编辑里把 wecom_external_user_id 填好。
WECOM_ROBOT_WEBHOOK=
WECOM_CORP_ID=
WECOM_AGENT_ID=
WECOM_SECRET=
WECOM_NOTIFY_DUE_DAYS=3
WECOM_NOTIFY_ENABLED=true
# 每天北京时间 09:00 自动执行采集计划。
ENABLE_SCHEDULER=true
SEED_DEMO_DATA=false

View File

@@ -494,6 +494,9 @@ export default function Home({ currentUser }: { currentUser: AuthUser }) {
const [activeNav, setActiveNav] = useState<NavKey>("overview");
const [loading, setLoading] = useState(true);
const [working, setWorking] = useState(false);
const [collectingDistributionIds, setCollectingDistributionIds] = useState<Set<string>>(
new Set(),
);
const [error, setError] = useState("");
const [toast, setToast] = useState("");
const [menuOpen, setMenuOpen] = useState(false);
@@ -762,10 +765,23 @@ export default function Home({ currentUser }: { currentUser: AuthUser }) {
};
const collectMetrics = async (distribution: Distribution) => {
setCollectingDistributionIds((current) => {
const next = new Set(current);
next.add(distribution.id);
return next;
});
try {
await runAction(
{ action: "collect_now", distributionId: distribution.id },
"公开数据已更新",
);
} finally {
setCollectingDistributionIds((current) => {
const next = new Set(current);
next.delete(distribution.id);
return next;
});
}
};
const updateDistributionPublishUrl = async (
@@ -1123,6 +1139,7 @@ export default function Home({ currentUser }: { currentUser: AuthUser }) {
tasks={data.tasks}
distributions={data.distributions}
working={working}
collectingDistributionIds={collectingDistributionIds}
canUpdatePublishUrl={isManager}
onCollect={collectMetrics}
onUpdatePublishUrl={updateDistributionPublishUrl}
@@ -2815,6 +2832,7 @@ function RecoveryPage({
tasks,
distributions,
working,
collectingDistributionIds,
canUpdatePublishUrl,
onCollect,
onUpdatePublishUrl,
@@ -2829,6 +2847,7 @@ function RecoveryPage({
tasks: Task[];
distributions: Distribution[];
working: boolean;
collectingDistributionIds: ReadonlySet<string>;
canUpdatePublishUrl: boolean;
onCollect: (distribution: Distribution) => void;
onUpdatePublishUrl: (distribution: Distribution) => void;
@@ -3116,7 +3135,9 @@ function RecoveryPage({
</span>
) : item.screenshot_key ? (
<span className="creator-ocr-failed">
KOC填写数据
{item.ocr_status === "processing"
? "正在识别数据"
: "识别失败,请手动填写"}
</span>
) : null}
{item.screenshot_key &&
@@ -3139,7 +3160,12 @@ function RecoveryPage({
) : (
<button
className="collect-button"
disabled={working || !noteUrl}
disabled={
!noteUrl ||
(working &&
(collectingDistributionIds.size === 0 ||
collectingDistributionIds.has(item.id)))
}
onClick={() => onCollect(item)}
>
{noteUrl ? "立即采集" : "待填链接"}

View File

@@ -37,6 +37,13 @@ import {
DistributionReleaseError,
releaseUnfinishedDistribution,
} from "../../../lib/distribution-release-service";
import {
resolveWecomConfig,
sendWecomAppMessage,
sendWecomRobotMessage,
WecomClientError,
type WecomBindings,
} from "../../../lib/wecom-client";
import { isManagerRequest } from "../../../lib/user-auth";
import { extractPublishUrl } from "../../../lib/publish-url";
@@ -534,6 +541,66 @@ export async function POST(request: Request) {
)
.bind(exposure, views, distributionId)
.run();
} else if (body.action === "bind_wecom_external_id") {
const partnerId = String(body.partnerId ?? "").trim().slice(0, 80);
const externalId = String(body.wecomExternalUserId ?? "")
.trim()
.slice(0, 128);
if (!partnerId) {
return Response.json(
{ error: "缺少 partnerId" },
{ status: 400 },
);
}
await db
.prepare(
"UPDATE partners SET wecom_external_user_id = ? WHERE id = ?",
)
.bind(externalId || null, partnerId)
.run();
return Response.json({
partnerId,
wecomExternalUserId: externalId || null,
});
} else if (body.action === "send_test_wecom") {
const wecomConfig = resolveWecomConfig(
env as unknown as WecomBindings,
);
const partnerId = String(body.partnerId ?? "").trim();
let partnerExternalId: string | null = null;
if (partnerId) {
const row = await db
.prepare(
"SELECT wecom_external_user_id FROM partners WHERE id = ?",
)
.bind(partnerId)
.first<{ wecom_external_user_id: string | null }>();
partnerExternalId = row?.wecom_external_user_id ?? null;
}
const testContent = `[KOC LOOP 测试] 群机器人连通性正常,时间 ${new Date().toISOString()}`;
let robotStatus: "ok" | "skipped" = "skipped";
if (wecomConfig.robotWebhook) {
await sendWecomRobotMessage(testContent, wecomConfig);
robotStatus = "ok";
}
let appStatus: "ok" | "skipped" | "failed" = "skipped";
if (
partnerExternalId &&
wecomConfig.corpId &&
wecomConfig.agentId &&
wecomConfig.secret
) {
const result = await sendWecomAppMessage(
[partnerExternalId],
testContent,
wecomConfig,
);
appStatus = result.failed > 0 ? "failed" : "ok";
}
return Response.json({
robot: robotStatus,
app: appStatus,
});
} else {
return Response.json({ error: "不支持的操作" }, { status: 400 });
}
@@ -545,7 +612,8 @@ export async function POST(request: Request) {
{
status:
error instanceof FeishuSourceError ||
error instanceof DistributionReleaseError
error instanceof DistributionReleaseError ||
error instanceof WecomClientError
? error.status
: 500,
},

View File

@@ -4,9 +4,10 @@ import {
getUploadBucket,
} from "../../../lib/mvp-db";
import { adminForbidden, isAdminRequest } from "../../../lib/admin-auth";
import { verifyCreatorScreenshotAccessToken } from "../../../lib/creator-screenshot-access";
export async function GET(request: Request) {
if (!(await isAdminRequest(request))) return adminForbidden();
const isAdmin = await isAdminRequest(request);
try {
await ensureSchema();
const distributionId = new URL(request.url).searchParams
@@ -22,6 +23,18 @@ export async function GET(request: Request) {
)
.bind(distributionId)
.first<{ screenshot_key: string | null }>();
const mcpToken = new URL(request.url).searchParams.get("mcp_token") || "";
if (
!isAdmin &&
(!row?.screenshot_key ||
!verifyCreatorScreenshotAccessToken(
distributionId,
row.screenshot_key,
mcpToken,
))
) {
return adminForbidden();
}
if (
!row?.screenshot_key ||
!row.screenshot_key.startsWith("creator-center/")

View File

@@ -13,6 +13,13 @@ import {
parseResultScreenshotKeys,
serializeResultScreenshotKeys,
} from "../../../lib/result-screenshots";
import { creatorScreenshotMcpUrl } from "../../../lib/creator-screenshot-access";
import {
extractCreatorMetricsFromMcp,
resolveCollectionMcpConfig,
type CollectionMcpBindings,
} from "../../../lib/mcp-collection-client";
import { getRuntimeEnv } from "../../../lib/runtime-env";
async function readUpload(request: Request) {
const contentType = request.headers.get("content-type") ?? "";
@@ -163,13 +170,47 @@ async function handlePost(request: Request) {
screenshot_key = ?,
ocr_status = CASE
WHEN exposure IS NOT NULL AND views IS NOT NULL THEN ocr_status
ELSE 'uploaded'
ELSE 'processing'
END,
updated_at = CURRENT_TIMESTAMP
WHERE id = ?`,
)
.bind(key, upload.distributionId)
.run();
try {
const metrics = await extractCreatorMetricsFromMcp(
creatorScreenshotMcpUrl(request, upload.distributionId, key),
resolveCollectionMcpConfig(
getRuntimeEnv() as unknown as CollectionMcpBindings,
),
);
await getRawDb()
.prepare(
`UPDATE distributions SET exposure = ?, views = ?,
ocr_status = 'success', updated_at = CURRENT_TIMESTAMP
WHERE id = ?`,
)
.bind(metrics.exposure, metrics.views, upload.distributionId)
.run();
return Response.json({
uploaded: true,
exposure: metrics.exposure,
views: metrics.views,
ocrStatus: "success",
kind: "creator-center",
});
} catch (error) {
console.error("[KOC LOOP] creator screenshot OCR failed", {
distributionId: upload.distributionId,
detail: error instanceof Error ? error.message : String(error),
});
await getRawDb()
.prepare(
`UPDATE distributions SET ocr_status = 'failed', updated_at = CURRENT_TIMESTAMP WHERE id = ?`,
)
.bind(upload.distributionId)
.run();
}
} else {
await getRawDb()
.prepare(
@@ -189,6 +230,7 @@ async function handlePost(request: Request) {
: isCreatorCenter
? "creator-center"
: "publish",
ocrStatus: isCreatorCenter ? "failed" : undefined,
});
} catch (error) {
return Response.json(

View File

@@ -6,6 +6,13 @@ import {
uid,
} from "../../../lib/mvp-db";
import { adminForbidden, isAdminRequest } from "../../../lib/admin-auth";
import { creatorScreenshotMcpUrl } from "../../../lib/creator-screenshot-access";
import {
extractCreatorMetricsFromMcp,
resolveCollectionMcpConfig,
type CollectionMcpBindings,
} from "../../../lib/mcp-collection-client";
import { getRuntimeEnv } from "../../../lib/runtime-env";
export async function POST(request: Request) {
if (!(await isAdminRequest(request))) return adminForbidden();
@@ -33,12 +40,37 @@ export async function POST(request: Request) {
screenshot_key = ?,
exposure = NULL,
views = NULL,
ocr_status = 'failed',
ocr_status = 'processing',
updated_at = CURRENT_TIMESTAMP
WHERE id = ?`,
)
.bind(key, distributionId)
.run();
try {
const metrics = await extractCreatorMetricsFromMcp(
creatorScreenshotMcpUrl(request, distributionId, key),
resolveCollectionMcpConfig(getRuntimeEnv() as unknown as CollectionMcpBindings),
);
await db
.prepare(
`UPDATE distributions SET exposure = ?, views = ?,
ocr_status = 'success', updated_at = CURRENT_TIMESTAMP
WHERE id = ?`,
)
.bind(metrics.exposure, metrics.views, distributionId)
.run();
} catch (error) {
console.error("[KOC LOOP] creator screenshot OCR failed", {
distributionId,
detail: error instanceof Error ? error.message : String(error),
});
await db
.prepare(
`UPDATE distributions SET ocr_status = 'failed', updated_at = CURRENT_TIMESTAMP WHERE id = ?`,
)
.bind(distributionId)
.run();
}
return Response.json(await getDashboardData());
} catch (error) {
return Response.json(

View File

@@ -954,7 +954,7 @@ export default function Home() {
selectedItem: Assignment,
file: File,
kind: "publish" | "creator-center" | "task-result",
) => {
): Promise<{ exposure?: number | null; views?: number | null; ocrStatus?: string }> => {
const compressed = await compressScreenshot(file);
const headers: Record<string, string> = {
"Content-Type": compressed.type || "image/jpeg",
@@ -973,7 +973,12 @@ export default function Home() {
headers,
body: compressed,
});
const uploadResult = (await uploadResponse.json()) as { error?: string };
const uploadResult = (await uploadResponse.json()) as {
error?: string;
exposure?: number | null;
views?: number | null;
ocrStatus?: string;
};
if (!uploadResponse.ok) {
throw new Error(
uploadResult.error ||
@@ -984,6 +989,7 @@ export default function Home() {
: "发布截图上传失败"),
);
}
return uploadResult;
};
const submitNote = async (event: FormEvent) => {
@@ -1082,22 +1088,31 @@ export default function Home() {
setToast("请选择创作者中心截图");
return;
}
if (
!/^\d{1,12}$/.test(creatorExposure) ||
!/^\d{1,12}$/.test(creatorViews)
) {
setToast("请填写正确的曝光量和阅读量");
return;
}
try {
setCreatorWorking(true);
let exposureValue = creatorExposure;
let viewsValue = creatorViews;
if (creatorScreenshot) {
setCreatorStage("正在上传截图…");
await uploadEvidence(
const ocrResult = await uploadEvidence(
selected,
creatorScreenshot,
"creator-center",
);
if (typeof ocrResult.exposure === "number") {
exposureValue = String(ocrResult.exposure);
setCreatorExposure(exposureValue);
}
if (typeof ocrResult.views === "number") {
viewsValue = String(ocrResult.views);
setCreatorViews(viewsValue);
}
if (ocrResult.ocrStatus === "success") {
setCreatorStage("已识别截图数据,正在保存…");
}
}
if (!/^\d{1,12}$/.test(exposureValue) || !/^\d{1,12}$/.test(viewsValue)) {
throw new Error("未识别到完整数据,请填写曝光量和阅读量");
}
setCreatorStage("正在保存数据…");
const response = await fetch(partnerApi("/api/partner"), {
@@ -1109,8 +1124,8 @@ export default function Home() {
claimToken,
delegationToken,
distributionId: selected.id,
exposure: creatorExposure,
views: creatorViews,
exposure: exposureValue,
views: viewsValue,
}),
});
const result = (await response.json()) as { error?: string };
@@ -1593,7 +1608,7 @@ export default function Home() {
<div className="evidence-empty-state">
<b></b>
<strong>7</strong>
<small>OCR</small>
<small></small>
</div>
)}
<label className="evidence-upload-action">

View File

@@ -92,7 +92,7 @@ test("keeps claiming minimal and backfill one-to-one", async () => {
assert.match(page, /submit_creator_metrics/);
assert.match(page, /creatorExposure/);
assert.match(page, /creatorViews/);
assert.match(page, /截图仅用于运营核对不再自动OCR/);
assert.match(page, /上传后自动识别曝光量和阅读量,识别失败可手动填写/);
assert.match(page, /evidenceImageUrl/);
assert.match(page, /params\.set\("v", evidenceKey\)/);
assert.match(page, /evidenceImageUrl\(selected, "publish"\)/);
@@ -140,7 +140,7 @@ test("renders screenshot-only tasks with anonymous delegation and multi-image up
assert.match(styles, /\.evidence-preview-button/);
});
test("shows D1 timestamps in Beijing time", () => {
test("shows D1 and MySQL UTC timestamps in Beijing time", () => {
const stored = "2026-07-29 05:36:00";
assert.equal(parseStoredDate(stored).toISOString(), "2026-07-29T05:36:00.000Z");
assert.match(formatShanghaiDate(stored, true), /07\/29.*13:36/);

View File

@@ -21,6 +21,33 @@ type ScheduledTask = {
type CollectionSource = "automatic" | "catchup" | "manual";
function userFacingCollectionError(error: unknown) {
const message = error instanceof Error ? error.message : String(error ?? "");
const normalized = message.toLowerCase();
if (
/cookie|登录|登陆|授权|access.?token|未登录|未授权|401|403/.test(
normalized,
)
) {
return "采集 Cookie 已过期或无权限";
}
if (
/链接|link|url|404|not found|不存在|删除|失效|无法识别.*笔记/.test(
normalized,
)
) {
return "笔记链接失效或不可访问";
}
if (
/超时|timeout|fetch failed|network|502|503|暂时不可用|响应空/.test(
normalized,
)
) {
return "采集服务暂时不可用";
}
return "笔记链接失效或采集 Cookie 过期";
}
function utcDay(value: string) {
const match = value.match(/^(\d{4})-(\d{2})-(\d{2})$/);
if (!match) return null;
@@ -311,8 +338,9 @@ export async function collectDistributionMetrics(
]);
return { skipped: false, likes, comments, collects, shares };
} catch (error) {
const message =
error instanceof Error ? error.message : "公开数据采集失败";
const detail = error instanceof Error ? error.message : String(error);
const message = userFacingCollectionError(error);
console.error("[KOC LOOP] collection failed", { detail, distributionId });
await db.batch([
db
.prepare(

View File

@@ -0,0 +1,60 @@
import { createHmac, timingSafeEqual } from "node:crypto";
import { getRuntimeEnv } from "./runtime-env";
function accessSecret() {
const env = getRuntimeEnv();
return (
env.KOC_MCP_API_KEY ||
env.KOC_LOOP_MCP_API_KEY ||
env.ADMIN_INTERNAL_TOKEN ||
env.AI_TOOL_CENTER_MCP_KEY ||
""
).trim();
}
function signature(distributionId: string, screenshotKey: string, expiresAt: number) {
return createHmac("sha256", accessSecret())
.update(`${distributionId}:${screenshotKey}:${expiresAt}`)
.digest("hex");
}
export function creatorScreenshotAccessToken(
distributionId: string,
screenshotKey: string,
expiresAt = Math.floor(Date.now() / 1000) + 300,
) {
const secret = accessSecret();
if (!secret) throw new Error("图片识别访问密钥未配置");
return `${expiresAt}.${signature(distributionId, screenshotKey, expiresAt)}`;
}
export function verifyCreatorScreenshotAccessToken(
distributionId: string,
screenshotKey: string,
token: string,
) {
const [expiresText, received] = token.split(".");
const expiresAt = Number(expiresText);
if (!Number.isInteger(expiresAt) || expiresAt < Math.floor(Date.now() / 1000)) {
return false;
}
const expected = signature(distributionId, screenshotKey, expiresAt);
if (!received || received.length !== expected.length) return false;
return timingSafeEqual(Buffer.from(received), Buffer.from(expected));
}
export function creatorScreenshotMcpUrl(
request: Request,
distributionId: string,
screenshotKey: string,
) {
const env = getRuntimeEnv();
const origin = (env.APP_ORIGIN || new URL(request.url).origin).replace(/\/$/, "");
const url = new URL(`${origin}/api/creator-screenshot`);
url.searchParams.set("distribution", distributionId);
url.searchParams.set(
"mcp_token",
creatorScreenshotAccessToken(distributionId, screenshotKey),
);
return url.toString();
}

View File

@@ -254,6 +254,7 @@ async function invokeMcpTool(
timeoutMs: number,
name: string,
args: Record<string, unknown>,
allowText = false,
): Promise<ToolResult> {
const result = await postMcp(
fetchImpl,
@@ -277,6 +278,12 @@ async function invokeMcpTool(
try {
payload = JSON.parse(text);
} catch {
if (allowText) {
return {
isError: result.envelope?.result?.isError === true,
payload: { text },
};
}
if (result.envelope?.result?.isError === true) {
return {
isError: true,
@@ -305,6 +312,7 @@ async function callMcpTool(
timeoutMs: number,
name: string,
args: Record<string, unknown>,
allowText = false,
): Promise<ToolResult> {
const nested = await invokeMcpTool(
fetchImpl,
@@ -313,6 +321,7 @@ async function callMcpTool(
timeoutMs,
name,
{ request: args },
allowText,
);
if (!isToolArgumentShapeError(nested)) return nested;
return invokeMcpTool(
@@ -322,10 +331,86 @@ async function callMcpTool(
timeoutMs,
name,
args,
allowText,
);
}
function metricsFromToolResult(result: ToolResult): XhsPublicMetrics {
function imageToolText(value: unknown) {
if (typeof value === "string") return value;
const root = asRecord(value);
if (!root) return "";
return [root.text, root.description, root.content, root.message, root.error]
.flatMap((item) => (Array.isArray(item) ? item : [item]))
.map((item) => {
if (typeof item === "string") return item;
const record = asRecord(item);
return String(record?.text ?? record?.description ?? "");
})
.filter(Boolean)
.join("\n");
}
function metricFromImageText(text: string, labels: string[], label: string) {
const pattern = labels.join("|");
const match = text.match(
new RegExp(`(?:${pattern})\\s*[:]?\\s*([\\d,.]+(?:万|w|千|k)?)`, "i"),
);
if (!match) return null;
try {
return metricValue(match[1], label);
} catch {
return null;
}
}
function creatorMetricsFromImageResult(result: ToolResult) {
if (result.isError) {
const detail = imageToolText(result.payload);
throw new Error(
`图片识别 MCP 调用失败${detail ? `${safeMessage(detail, "")}` : ""}`,
);
}
const text = imageToolText(result.payload);
const exposure = metricFromImageText(
text,
["曝光量", "曝光", "impressions", "exposure"],
"曝光量",
);
const views = metricFromImageText(
text,
["阅读量", "阅读", "views", "view_count", "view count"],
"阅读量",
);
if (exposure === null || views === null) {
throw new Error("图片识别未找到曝光量和阅读量");
}
return { exposure, views };
}
export async function extractCreatorMetricsFromMcp(
imageUrl: string,
config: CollectionMcpConfig,
fetchImpl: typeof fetch = fetch,
) {
const endpoint = buildMcpUrl(config);
const timeoutMs = Math.max(5_000, config.timeoutMs ?? 30_000);
const sessionId = await createMcpSession(fetchImpl, endpoint, timeoutMs);
const result = await callMcpTool(
fetchImpl,
endpoint,
sessionId,
timeoutMs,
"analyze_image",
{ image_url: imageUrl },
true,
);
return creatorMetricsFromImageResult(result);
}
function metricsFromToolResult(
result: ToolResult,
toolName = "fetch_content_detail",
): XhsPublicMetrics {
const root = asRecord(result.payload);
const response = asRecord(root?.response) ?? root;
const data = asRecord(response?.data) ?? asRecord(root?.data);
@@ -337,14 +422,18 @@ function metricsFromToolResult(result: ToolResult): XhsPublicMetrics {
success === false ||
(Number.isFinite(code) && code >= 400)
) {
const providerCode = Number(response?.code);
const codeLabel = Number.isFinite(providerCode)
? `code ${providerCode}`
: "";
throw new Error(
safeMessage(
`MCP工具 ${toolName} 返回失败${codeLabel}${safeMessage(
response?.msg ?? response?.message ?? root?.message,
"公开数据采集失败",
),
)}`,
);
}
if (!data) throw new Error("采集结果缺少互动数据");
if (!data) throw new Error(`MCP工具 ${toolName} 未返回互动数据`);
const count = (value: unknown, label: string) =>
value === null || value === undefined || value === ""
@@ -1156,7 +1245,7 @@ async function collectInSession(
auto_cookie: true,
},
);
return metricsFromToolResult(primary);
return metricsFromToolResult(primary, "fetch_content_detail");
}
export function resolveCollectionMcpConfig(

View File

@@ -20,6 +20,12 @@ export type RuntimeEnv = {
AI_TOOL_CENTER_MCP_KEY?: string;
COLLECTION_MCP_URL?: string;
COLLECTION_MCP_KEY?: string;
WECOM_CORP_ID?: string;
WECOM_AGENT_ID?: string;
WECOM_SECRET?: string;
WECOM_ROBOT_WEBHOOK?: string;
WECOM_NOTIFY_DUE_DAYS?: string;
WECOM_NOTIFY_ENABLED?: string;
SEED_DEMO_DATA?: string;
ENABLE_SCHEDULER?: string;
};

View File

@@ -8,6 +8,11 @@ import {
} from "./mcp-collection-client";
import { ensureSchema, getRawDb } from "./mvp-db";
import { getRuntimeEnv, isEnabled } from "./runtime-env";
import {
resolveWecomConfig,
type WecomBindings,
} from "./wecom-client";
import { runDueSoonWecomNotifications } from "./wecom-notifier-service";
declare global {
var __kocLoopScheduler: ScheduledTask | undefined;
@@ -17,14 +22,31 @@ async function runDailyJob() {
await withDatabaseLock("koc-loop-daily-collection", 0, async () => {
await ensureSchema();
const db = getRawDb();
const env = getRuntimeEnv();
const config = resolveCollectionMcpConfig(
getRuntimeEnv() as unknown as CollectionMcpBindings,
env as unknown as CollectionMcpBindings,
);
const collections = await runScheduledCollections(db, Date.now(), config);
const accounts = await backfillAccountProfiles(db, config, 10);
let wecom: Awaited<ReturnType<typeof runDueSoonWecomNotifications>> | null =
null;
if (isEnabled(env.WECOM_NOTIFY_ENABLED, true)) {
const wecomConfig = resolveWecomConfig(env as unknown as WecomBindings);
if (
wecomConfig.robotWebhook ||
(wecomConfig.corpId && wecomConfig.agentId && wecomConfig.secret)
) {
try {
wecom = await runDueSoonWecomNotifications(db, wecomConfig);
} catch (error) {
console.error("[KOC LOOP] wecom notify failed", error);
}
}
}
console.info("[KOC LOOP] daily scheduler completed", {
collections,
accounts,
wecom,
});
});
}

219
lib/wecom-client.ts Normal file
View File

@@ -0,0 +1,219 @@
const WECOM_API_ORIGIN = "https://qyapi.weixin.qq.com";
const DEFAULT_DUE_DAYS = 3;
export type WecomBindings = {
WECOM_CORP_ID?: string;
WECOM_AGENT_ID?: string;
WECOM_SECRET?: string;
WECOM_ROBOT_WEBHOOK?: string;
WECOM_NOTIFY_DUE_DAYS?: string;
};
export type WecomConfig = {
corpId: string;
agentId: string;
secret: string;
robotWebhook: string;
dueDays: number;
};
type FetchLike = typeof fetch;
type WecomEnvelope = {
errcode?: number;
errmsg?: string;
access_token?: string;
expires_in?: number;
invaliduser?: string;
};
type CachedAccessToken = {
corpId: string;
secret: string;
token: string;
expiresAt: number;
};
let cachedAccessToken: CachedAccessToken | null = null;
export class WecomClientError extends Error {
status: number;
constructor(message: string, status = 502) {
super(message);
this.name = "WecomClientError";
this.status = status;
}
}
function bindingValue(value: unknown) {
return String(value ?? "").trim();
}
function safeMessage(value: unknown) {
return String(value ?? "").trim().slice(0, 240);
}
export function resolveWecomConfig(bindings: WecomBindings): WecomConfig {
return {
corpId: bindingValue(bindings.WECOM_CORP_ID),
agentId: bindingValue(bindings.WECOM_AGENT_ID),
secret: bindingValue(bindings.WECOM_SECRET),
robotWebhook: bindingValue(bindings.WECOM_ROBOT_WEBHOOK),
dueDays: parseDueDays(bindings.WECOM_NOTIFY_DUE_DAYS),
};
}
function parseDueDays(value: string | undefined) {
const parsed = Number(bindingValue(value));
if (!Number.isFinite(parsed) || parsed < 1) return DEFAULT_DUE_DAYS;
return Math.min(30, Math.floor(parsed));
}
function hasAppCredentials(config: WecomConfig) {
return Boolean(config.corpId && config.agentId && config.secret);
}
async function readEnvelope(
response: Response,
fallbackMessage: string,
): Promise<WecomEnvelope> {
const text = await response.text();
try {
return JSON.parse(text) as WecomEnvelope;
} catch {
throw new WecomClientError(
`${fallbackMessage}(企业微信返回了非 JSON 响应)`,
502,
);
}
}
function ensureOk(
payload: WecomEnvelope,
fallbackMessage: string,
) {
const code = Number(payload.errcode ?? 0);
if (code === 0) return;
const message = safeMessage(payload.errmsg) || fallbackMessage;
if (code === 40014 || code === 42001) {
throw new WecomClientError(`企业微信 access_token 无效:${message}`, 401);
}
throw new WecomClientError(`${fallbackMessage}${message}`, 502);
}
async function fetchAccessToken(
config: WecomConfig,
fetchImpl: FetchLike,
) {
if (
cachedAccessToken?.corpId === config.corpId &&
cachedAccessToken?.secret === config.secret &&
cachedAccessToken.expiresAt > Date.now() + 60_000
) {
return cachedAccessToken.token;
}
const url = new URL(`${WECOM_API_ORIGIN}/cgi-bin/gettoken`);
url.searchParams.set("corpid", config.corpId);
url.searchParams.set("corpsecret", config.secret);
const response = await fetchImpl(url.toString(), {
method: "GET",
signal: AbortSignal.timeout(12_000),
});
const payload = await readEnvelope(response, "获取企业微信 access_token 失败");
if (!response.ok) {
throw new WecomClientError(
`获取企业微信 access_token 失败HTTP ${response.status}`,
502,
);
}
ensureOk(payload, "获取企业微信 access_token 失败");
const token = bindingValue(payload.access_token);
if (!token) {
throw new WecomClientError("企业微信未返回有效 access_token", 502);
}
cachedAccessToken = {
corpId: config.corpId,
secret: config.secret,
token,
expiresAt:
Date.now() + Math.max(300, Number(payload.expires_in) || 7_200) * 1_000,
};
return token;
}
export async function sendWecomRobotMessage(
content: string,
config: WecomConfig,
fetchImpl: FetchLike = fetch,
): Promise<void> {
if (!config.robotWebhook) {
console.warn("[KOC LOOP] wecom robot webhook not configured, skipping");
return;
}
const response = await fetchImpl(config.robotWebhook, {
method: "POST",
headers: { "Content-Type": "application/json; charset=utf-8" },
body: JSON.stringify({
msgtype: "text",
text: { content },
}),
signal: AbortSignal.timeout(10_000),
});
const payload = await readEnvelope(response, "企业微信群机器人推送失败");
if (!response.ok) {
throw new WecomClientError(
`企业微信群机器人推送失败HTTP ${response.status}`,
502,
);
}
ensureOk(payload, "企业微信群机器人推送失败");
}
export async function sendWecomAppMessage(
externalUserIds: string[],
content: string,
config: WecomConfig,
fetchImpl: FetchLike = fetch,
): Promise<{ sent: number; failed: number; skipped: boolean }> {
const normalized = externalUserIds
.map((id) => bindingValue(id))
.filter((id) => id.length > 0);
if (normalized.length === 0) {
return { sent: 0, failed: 0, skipped: true };
}
if (!hasAppCredentials(config)) {
return { sent: 0, failed: 0, skipped: true };
}
const token = await fetchAccessToken(config, fetchImpl);
const url = `${WECOM_API_ORIGIN}/cgi-bin/message/send?access_token=${encodeURIComponent(token)}`;
const response = await fetchImpl(url, {
method: "POST",
headers: { "Content-Type": "application/json; charset=utf-8" },
body: JSON.stringify({
touser: normalized.join("|"),
msgtype: "text",
agentid: Number(config.agentId),
text: { content },
}),
signal: AbortSignal.timeout(12_000),
});
const payload = await readEnvelope(response, "企业微信应用消息推送失败");
if (!response.ok) {
throw new WecomClientError(
`企业微信应用消息推送失败HTTP ${response.status}`,
502,
);
}
ensureOk(payload, "企业微信应用消息推送失败");
const invalid = bindingValue(payload.invaliduser).split("|").filter(Boolean);
return {
sent: Math.max(0, normalized.length - invalid.length),
failed: invalid.length,
skipped: false,
};
}
export function clearWecomAccessTokenCacheForTests() {
cachedAccessToken = null;
}

View File

@@ -0,0 +1,171 @@
import type { DatabaseClient } from "./database";
import { shanghaiDateFromTimestamp } from "./collection-service";
import {
sendWecomAppMessage,
sendWecomRobotMessage,
type WecomConfig,
} from "./wecom-client";
import { getRuntimeEnv } from "./runtime-env";
type FetchLike = typeof fetch;
export type WecomNotifySummary = {
dueSoonAttempted: number;
dueSoonSent: number;
dueSoonFailed: number;
dueSoonSkipped: number;
digestSent: boolean;
};
type DueSoonRow = {
distribution_id: string;
partner_id: string;
partner_name: string;
wecom_external_user_id: string | null;
task_name: string;
due_at: string;
content_title: string;
};
export function computeDueCutoff(today: string, dueDays: number) {
const match = /^(\d{4})-(\d{2})-(\d{2})$/.exec(today);
if (!match) return today;
const [, y, m, d] = match;
const date = new Date(
Date.UTC(Number(y), Number(m) - 1, Number(d)) + dueDays * 24 * 60 * 60 * 1_000,
);
return date.toISOString().slice(0, 10);
}
export async function runDueSoonWecomNotifications(
db: DatabaseClient,
config: WecomConfig,
fetchImpl: FetchLike = fetch,
now: number = Date.now(),
): Promise<WecomNotifySummary> {
const today = shanghaiDateFromTimestamp(now);
const cutoff = computeDueCutoff(today, config.dueDays);
const portalUrl = bindingValue(getRuntimeEnv().KOC_PORTAL_URL);
const result = await db
.prepare(
`SELECT
d.id AS distribution_id,
d.partner_id,
p.name AS partner_name,
p.wecom_external_user_id,
t.name AS task_name,
t.due_at,
c.title AS content_title
FROM distributions d
JOIN tasks t ON t.id = d.task_id
JOIN partners p ON p.id = d.partner_id
JOIN contents c ON c.id = d.content_id
WHERE (d.publish_url IS NULL OR d.publish_url = '')
AND t.due_at IS NOT NULL AND t.due_at != ''
AND t.due_at <= ?
ORDER BY t.due_at ASC, p.name ASC`,
)
.bind(cutoff)
.all<DueSoonRow>();
const grouped = new Map<
string,
{
partnerName: string;
externalUserId: string | null;
rows: DueSoonRow[];
}
>();
for (const row of result.results) {
const entry = grouped.get(row.partner_id) ?? {
partnerName: row.partner_name,
externalUserId: row.wecom_external_user_id,
rows: [],
};
entry.rows.push(row);
if (!entry.externalUserId && row.wecom_external_user_id) {
entry.externalUserId = row.wecom_external_user_id;
}
grouped.set(row.partner_id, entry);
}
let dueSoonAttempted = 0;
let dueSoonSent = 0;
let dueSoonFailed = 0;
let dueSoonSkipped = 0;
const digestTasks: string[] = [];
for (const [, entry] of grouped) {
dueSoonAttempted += 1;
const external = entry.externalUserId
? [entry.externalUserId]
: [];
const lines = entry.rows.slice(0, 5).map((row) => {
return `· 《${truncate(row.task_name, 24)}》— ${truncate(row.content_title, 24)}(截止 ${row.due_at}`;
});
const overflow =
entry.rows.length > 5 ? `\n…还有 ${entry.rows.length - 5}` : "";
const link = portalUrl ? `\n领取链接${portalUrl}` : "";
const content =
`${entry.partnerName},你有 ${entry.rows.length} 条内容待发布:\n${lines.join("\n")}${overflow}${link}`;
const appResult = await sendWecomAppMessage(
external,
content,
config,
fetchImpl,
).catch((error: unknown) => {
console.warn(
"[KOC LOOP] wecom app message failed",
{ partner: entry.partnerName, error: safeError(error) },
);
return null;
});
if (appResult === null) {
dueSoonFailed += 1;
} else if (appResult.skipped) {
dueSoonSkipped += 1;
} else {
dueSoonSent += 1;
}
const earliestDue = entry.rows[0]?.due_at ?? "";
digestTasks.push(
`· ${entry.partnerName}${entry.rows.length} 条,最近截止 ${earliestDue}`,
);
}
let digestSent = false;
if (config.robotWebhook && digestTasks.length > 0) {
const digest =
`今日待发布催办(${today},截止 ≤ ${cutoff}\n${digestTasks.join("\n")}`;
try {
await sendWecomRobotMessage(digest, config, fetchImpl);
digestSent = true;
} catch (error) {
console.warn(
"[KOC LOOP] wecom robot digest failed",
{ error: safeError(error) },
);
}
}
return {
dueSoonAttempted,
dueSoonSent,
dueSoonFailed,
dueSoonSkipped,
digestSent,
};
}
function bindingValue(value: unknown) {
return String(value ?? "").trim();
}
function truncate(value: string, max: number) {
return value.length > max ? `${value.slice(0, max)}` : value;
}
function safeError(error: unknown) {
return error instanceof Error ? error.message : String(error);
}

View File

@@ -6,7 +6,7 @@ import {
parseStoredDate,
} from "../lib/date-utils.ts";
test("treats D1 CURRENT_TIMESTAMP values as UTC and displays Beijing time", () => {
test("treats D1 and MySQL UTC timestamps as UTC and displays Beijing time", () => {
const stored = "2026-07-29 05:36:00";
assert.equal(parseStoredDate(stored).toISOString(), "2026-07-29T05:36:00.000Z");
assert.match(formatShanghaiDate(stored, true), /07\/29.*13:36/);

View File

@@ -520,7 +520,7 @@ test("surfaces current detail tool failures without calling removed tools", asyn
},
fetchImpl,
),
/获取内容详情失败/,
/MCP工具 fetch_content_detail 返回失败code 400获取内容详情失败/,
);
assert.equal(calls.length, 3);
assert.equal(calls[2].body.params.name, "fetch_content_detail");

View File

@@ -190,8 +190,9 @@ test("issues external task links and supports one-to-one note submissions", asyn
assert.doesNotMatch(partnerUtils, /user\/profile\/\$\{platformUid\}/);
assert.match(uploadRoute, /publish-evidence/);
assert.match(uploadRoute, /creator-center/);
assert.match(uploadRoute, /ELSE 'uploaded'/);
assert.doesNotMatch(uploadRoute, /ocrMetric/);
assert.match(uploadRoute, /ELSE 'processing'/);
assert.match(uploadRoute, /extractCreatorMetricsFromMcp/);
assert.match(uploadRoute, /creatorScreenshotMcpUrl/);
assert.match(uploadRoute, /x-koc-upload-kind/);
assert.match(uploadRoute, /x-koc-distribution/);
assert.match(uploadRoute, /request\.arrayBuffer/);
@@ -223,7 +224,7 @@ test("issues external task links and supports one-to-one note submissions", asyn
assert.match(adminApp, /全部内容类型/);
assert.match(adminApp, /全部平台/);
assert.match(adminApp, /task-scope-subline/);
assert.match(adminApp, /待KOC填写数据/);
assert.match(adminApp, /识别失败,请手动填写/);
assert.match(adminApp, /AdminImageLightbox/);
assert.match(adminApp, /CreatorScreenshotPreview/);
assert.match(adminApp, /admin-image-lightbox/);

217
tests/wecom-client.test.mjs Normal file
View File

@@ -0,0 +1,217 @@
import assert from "node:assert/strict";
import test from "node:test";
import {
clearWecomAccessTokenCacheForTests,
resolveWecomConfig,
sendWecomAppMessage,
sendWecomRobotMessage,
WecomClientError,
} from "../lib/wecom-client.ts";
const baseBindings = {
WECOM_CORP_ID: "corp-test",
WECOM_AGENT_ID: "1000001",
WECOM_SECRET: "secret-test",
WECOM_ROBOT_WEBHOOK:
"https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=robot-test",
WECOM_NOTIFY_DUE_DAYS: "3",
};
function fakeWecom(options = {}) {
const calls = [];
const fetchImpl = async (input, init = {}) => {
const url = new URL(String(input));
calls.push({ url, init });
if (url.pathname.endsWith("/cgi-bin/gettoken")) {
return Response.json({
errcode: 0,
errmsg: "ok",
access_token: "token-test",
expires_in: 7200,
});
}
if (url.pathname.endsWith("/cgi-bin/webhook/send")) {
return Response.json({
errcode: options.robotErrcode ?? 0,
errmsg: options.robotErrmsg ?? "ok",
});
}
if (url.pathname.endsWith("/cgi-bin/message/send")) {
const body = JSON.parse(String(init.body ?? "{}"));
const invalid = options.invalidUser
? body.touser.split("|").filter((u) => u === options.invalidUser)
: [];
return Response.json({
errcode: 0,
errmsg: "ok",
invaliduser: invalid.join("|"),
});
}
return new Response("not found", { status: 404 });
};
return { calls, fetchImpl };
}
test.beforeEach(() => {
clearWecomAccessTokenCacheForTests();
});
test("resolveWecomConfig applies dueDays and trims values", () => {
const config = resolveWecomConfig({
WECOM_CORP_ID: " corp ",
WECOM_AGENT_ID: " 10 ",
WECOM_SECRET: " s ",
WECOM_ROBOT_WEBHOOK: " https://hook ",
WECOM_NOTIFY_DUE_DAYS: "5",
});
assert.equal(config.corpId, "corp");
assert.equal(config.agentId, "10");
assert.equal(config.secret, "s");
assert.equal(config.robotWebhook, "https://hook");
assert.equal(config.dueDays, 5);
});
test("resolveWecomConfig falls back to default dueDays for bad input", () => {
assert.equal(resolveWecomConfig({ WECOM_NOTIFY_DUE_DAYS: "0" }).dueDays, 3);
assert.equal(resolveWecomConfig({ WECOM_NOTIFY_DUE_DAYS: "abc" }).dueDays, 3);
assert.equal(
resolveWecomConfig({ WECOM_NOTIFY_DUE_DAYS: "100" }).dueDays,
30,
);
});
test("sendWecomRobotMessage posts text payload and returns void", async () => {
const { calls, fetchImpl } = fakeWecom();
await sendWecomRobotMessage(
"hello",
resolveWecomConfig(baseBindings),
fetchImpl,
);
const hookCall = calls.find((c) =>
c.url.pathname.endsWith("/cgi-bin/webhook/send"),
);
assert.ok(hookCall, "robot webhook was called");
const body = JSON.parse(String(hookCall.init.body));
assert.equal(body.msgtype, "text");
assert.equal(body.text.content, "hello");
});
test("sendWecomRobotMessage skips silently when webhook is empty", async () => {
const { calls, fetchImpl } = fakeWecom();
const config = resolveWecomConfig({ ...baseBindings, WECOM_ROBOT_WEBHOOK: "" });
await sendWecomRobotMessage("hello", config, fetchImpl);
assert.equal(
calls.filter((c) =>
c.url.pathname.endsWith("/cgi-bin/webhook/send"),
).length,
0,
);
});
test("sendWecomRobotMessage throws WecomClientError on provider error", async () => {
const { fetchImpl } = fakeWecom({
robotErrcode: 93000,
robotErrmsg: "invalid webhook url",
});
await assert.rejects(
() =>
sendWecomRobotMessage(
"hello",
resolveWecomConfig(baseBindings),
fetchImpl,
),
(error) => {
assert.ok(error instanceof WecomClientError);
assert.equal(error.status, 502);
return true;
},
);
});
test("sendWecomAppMessage fetches access_token then sends to message/send", async () => {
const { calls, fetchImpl } = fakeWecom();
const result = await sendWecomAppMessage(
["user-a", "user-b"],
"催办内容",
resolveWecomConfig(baseBindings),
fetchImpl,
);
assert.equal(result.sent, 2);
assert.equal(result.failed, 0);
assert.equal(result.skipped, false);
const tokenCall = calls.find((c) =>
c.url.pathname.endsWith("/cgi-bin/gettoken"),
);
assert.ok(tokenCall, "gettoken was called");
const sendCall = calls.find((c) =>
c.url.pathname.endsWith("/cgi-bin/message/send"),
);
assert.ok(sendCall, "message/send was called");
assert.equal(sendCall.url.searchParams.get("access_token"), "token-test");
const body = JSON.parse(String(sendCall.init.body));
assert.equal(body.touser, "user-a|user-b");
assert.equal(body.agentid, 1000001);
assert.equal(body.text.content, "催办内容");
});
test("sendWecomAppMessage skips when no external user ids", async () => {
const { calls, fetchImpl } = fakeWecom();
const result = await sendWecomAppMessage(
[],
"催办",
resolveWecomConfig(baseBindings),
fetchImpl,
);
assert.equal(result.skipped, true);
assert.equal(
calls.filter((c) =>
c.url.pathname.endsWith("/cgi-bin/message/send"),
).length,
0,
);
});
test("sendWecomAppMessage skips when corp credentials are missing", async () => {
const { calls, fetchImpl } = fakeWecom();
const result = await sendWecomAppMessage(
["user-a"],
"催办",
resolveWecomConfig({
...baseBindings,
WECOM_CORP_ID: "",
WECOM_AGENT_ID: "",
WECOM_SECRET: "",
}),
fetchImpl,
);
assert.equal(result.skipped, true);
assert.equal(
calls.filter((c) =>
c.url.pathname.endsWith("/cgi-bin/gettoken"),
).length,
0,
);
});
test("sendWecomAppMessage reports failed when provider marks invaliduser", async () => {
const { fetchImpl } = fakeWecom({ invalidUser: "user-a" });
const result = await sendWecomAppMessage(
["user-a", "user-b"],
"催办",
resolveWecomConfig(baseBindings),
fetchImpl,
);
assert.equal(result.sent, 1);
assert.equal(result.failed, 1);
});
test("access_token is cached across calls", async () => {
const { calls, fetchImpl } = fakeWecom();
const config = resolveWecomConfig(baseBindings);
await sendWecomAppMessage(["u1"], "a", config, fetchImpl);
await sendWecomAppMessage(["u2"], "b", config, fetchImpl);
const tokenCalls = calls.filter((c) =>
c.url.pathname.endsWith("/cgi-bin/gettoken"),
);
assert.equal(tokenCalls.length, 1, "access_token cached for second call");
});

View File

@@ -0,0 +1,15 @@
import assert from "node:assert/strict";
import test from "node:test";
import { computeDueCutoff } from "../lib/wecom-notifier-service.ts";
test("computeDueCutoff advances by dueDays and crosses month/year", () => {
assert.equal(computeDueCutoff("2026-08-18", 3), "2026-08-21");
assert.equal(computeDueCutoff("2026-08-30", 3), "2026-09-02");
assert.equal(computeDueCutoff("2026-12-30", 3), "2027-01-02");
assert.equal(computeDueCutoff("2026-08-18", 0), "2026-08-18");
assert.equal(computeDueCutoff("2026-08-18", 10), "2026-08-28");
});
test("computeDueCutoff returns input unchanged when malformed", () => {
assert.equal(computeDueCutoff("not-a-date", 3), "not-a-date");
});