feat(hub): 飞书 @bot 触发链路 + ProjectAgentLock 不变式(ADR-0001/0002)

群消息事件 → 创建 AgentRun → 获锁 → 调 runAgent → 释放锁 + 状态回写:
- feishu/client.ts: lark WSClient 长连接接收 im.message.receive_v1;
  sendText 经 client.im.v1 发送(SDK 运行时动态 API,类型弱化务实收口)
- feishu/trigger.ts: @bot 识别(chat 绑定解析 ADR-0001)→ 创建 run+
  session(provider-bound ADR-0017)→ acquireLock(ADR-0002,竞态回退 busy)→
  runAgent(fire-and-forget)→ finally releaseLock;非项目群/无 @bot/非文本
  全部忽略
- lock.ts: ADR-0002 WellFormed 落代码——currentLockRunId 读锁时校验
  持锁 run 非终止,终止 run 持锁即 release-path bug(抛而非静默)
- models.ts: resolve 提到 ModelRegistry 接口(之前漏在实现上)
- server.ts: Fastify + lark ws 启动,/api/healthz
prompt 提取 smoke 通过(@bot+prompt→prompt;纯@bot→null;无mention→null;
非文本→null)。修了 mention 正则 bug(/@_\w+_/ 的回溯错误)。tsc rc=0。
This commit is contained in:
2026-07-06 23:11:35 +08:00
parent b4dd8f09e1
commit 313e890b6c
7 changed files with 772 additions and 42 deletions
+2
View File
@@ -39,6 +39,8 @@ export interface ModelRegistry {
listRoles(): readonly RoleEntry[];
/** The role preset with the given id, if any. */
role(id: string): RoleEntry | undefined;
/** Resolve a model id: explicit request → role default → first enabled. */
resolve(requestedModel: string | undefined, roleId: string): string;
}
/**
+102
View File
@@ -0,0 +1,102 @@
/**
* Feishu lark client + long-connection event dispatcher.
*
* Uses lark's WebSocket long-connection mode (no public callback endpoint). The
* `EventDispatcher` registers `im.message.receive_v1`; the SDK delivers the raw
* event to our handler. This version of the SDK does not auto-normalize on the
* ws path — `normalize()` is a separate helper — so we work with the raw event
* shape directly (it carries chat_id, message_type, content, mentions).
*
* ADR-0001: the bot operates inside a project's bound group. It does not own
* the lock (ADR-0002) and is not the session owner.
*/
import * as lark from "@larksuiteoapi/node-sdk";
import type { FastifyBaseLogger } from "fastify";
export interface FeishuConfig {
readonly appId: string;
readonly appSecret: string;
readonly botOpenId: string;
}
export interface FeishuRuntime {
readonly client: lark.Client;
readonly logger: FastifyBaseLogger;
}
/** Raw `im.message.receive_v1` event payload (subset we use). */
export interface MessageReceiveEvent {
readonly message: {
readonly message_id: string;
readonly chat_id: string;
readonly chat_type: string;
readonly message_type: string;
readonly content: string;
readonly mentions?: ReadonlyArray<{
readonly key: string;
readonly id: { readonly open_id?: string };
readonly name: string;
}>;
};
readonly sender: {
readonly sender_id: { readonly open_id?: string };
readonly sender_type: string;
};
}
/** Send a text message to a chat. */
export async function sendText(rt: FeishuRuntime, chatId: string, text: string): Promise<void> {
// SDK exposes im.v1 dynamically at runtime; the generated types are weak
// there, so cast through the known shape.
const client = rt.client as unknown as {
im: { v1: { message: { create: (p: unknown) => Promise<unknown> } } };
};
await client.im.v1.message.create({
params: { receive_id_type: "chat_id" },
data: {
receive_id: chatId,
msg_type: "text",
content: JSON.stringify({ text }),
},
});
}
/**
* Start the lark WebSocket client and route `im.message.receive_v1` events to
* `onMessage`. The ws connection lives in the background; this resolves the
* runtime handle.
*/
export function startFeishuListener(
config: FeishuConfig,
logger: FastifyBaseLogger,
onMessage: (event: MessageReceiveEvent, rt: FeishuRuntime) => Promise<void>,
): FeishuRuntime {
const client = new lark.Client({
appId: config.appId,
appSecret: config.appSecret,
appType: lark.AppType.SelfBuild,
});
const wsClient = new lark.WSClient({
appId: config.appId,
appSecret: config.appSecret,
domain: lark.Domain.Feishu,
loggerLevel: lark.LoggerLevel.info,
});
const rt: FeishuRuntime = { client, logger };
void wsClient.start({
eventDispatcher: new lark.EventDispatcher({}).register({
"im.message.receive_v1": async (data) => {
try {
await onMessage(data as MessageReceiveEvent, rt);
} catch (e) {
logger.error({ err: e }, "feishu message handler threw");
}
},
}),
});
return rt;
}
+211
View File
@@ -0,0 +1,211 @@
/**
* Feishu @bot trigger → AgentRun lifecycle.
*
* The core path (ADR-0001..0003, 0017):
* 1. A message mentioning the bot arrives in a project group.
* 2. Resolve the chat → project binding (ADR-0001). Unknown chat ⇒ ignored.
* 3. Create an AgentRun + AgentSession (provider-bound, ADR-0017).
* 4. Acquire the project lock for the run (ADR-0002). If locked, reply busy.
* 5. Run the agent loop; the agent reads Feishu context on demand through the
* ADR-0003-authorized tool (chat must match binding).
* 6. On finish/fail/timeout: release the lock, post a status card.
*
* The agent provider/model come from the ModelRegistry (ADR-0017 role-based
* routing, role-as-data). The trigger does not pick a model directly.
*/
import type { PrismaClient } from "@prisma/client";
import type { FastifyBaseLogger } from "fastify";
import { sendText, type FeishuRuntime, type MessageReceiveEvent } from "./client.js";
import type { AgentProvider } from "../agent/provider.js";
import type { ToolRegistry } from "../agent/tools.js";
import type { ModelRegistry } from "../agent/models.js";
import { runAgent, type ProjectContext } from "../agent/runner.js";
import { acquireLock, currentLockRunId, releaseLock } from "../lock.js";
interface TriggerDeps {
readonly prisma: PrismaClient;
readonly provider: AgentProvider;
readonly tools: ToolRegistry;
readonly models: ModelRegistry;
readonly logger: FastifyBaseLogger;
}
/**
* Build the message handler. Returns a function suitable for the lark
* `im.message.receive_v1` event. The handler is async and long-running (the
* agent run is the long leg); the lock guarantees one run per project at a time.
*/
export function makeTriggerHandler(deps: TriggerDeps) {
return async (event: MessageReceiveEvent, rt: FeishuRuntime): Promise<void> => {
const msg = event.message;
const sender = event.sender;
void sender; // sender→user mapping is OPEN (principal sub-typology)
// Only react to text messages in group chats.
if (msg.message_type !== "text") return;
const chatId = msg.chat_id;
if (chatId === "") return;
// ADR-0001: resolve chat → project. Unknown chat ⇒ not a project group.
const binding = await deps.prisma.projectGroupBinding.findUnique({
where: { chatId },
select: { projectId: true },
});
if (binding === null) {
deps.logger.debug({ chatId }, "feishu trigger: chat not bound to any project");
return;
}
const projectId = binding.projectId;
// ADR-0002: if locked by another run, reply busy.
const existing = await currentLockRunId(deps.prisma, projectId);
if (existing !== null) {
await sendText(rt, chatId, "项目正在处理中,请稍候。");
return;
}
const prompt = extractPrompt(msg);
if (prompt === null) return;
const project = await deps.prisma.project.findUnique({
where: { id: projectId },
select: { workspaceDir: true },
});
if (project === null) {
deps.logger.error({ projectId }, "feishu trigger: project missing");
return;
}
// ADR-0017: role-as-data; model resolved by registry. "draft" is the
// fallback role id; real wiring lets the trigger carry a role hint (slash
// command, card selector). The registry degrades unknown roles gracefully.
const model = deps.models.resolve(undefined, "draft");
const providerId = deps.provider.id;
// ADR-0017: provider+model-bound session. Reuse if one exists for this
// (project, provider, model) triple; otherwise create.
const existingSession = await deps.prisma.agentSession.findFirst({
where: { projectId, provider: providerId, model, archivedAt: null },
select: { id: true },
});
const session =
existingSession !== null
? await deps.prisma.agentSession.update({
where: { id: existingSession.id },
data: { provider: providerId, model, updatedAt: new Date() },
select: { id: true },
})
: await deps.prisma.agentSession.create({
data: {
projectId,
provider: providerId,
model,
title: prompt.slice(0, 40) || null,
metadata: {},
},
select: { id: true },
});
const run = await deps.prisma.agentRun.create({
data: {
projectId,
sessionId: session.id,
requestedByUserId: null, // sender→user mapping is OPEN (principal sub-typology)
entrypoint: "FEISHU",
status: "ACTIVE",
prompt,
model,
provider: providerId,
metadata: {},
},
select: { id: true },
});
// ADR-0002: acquire the project lock for this run.
try {
await acquireLock(deps.prisma, projectId, run.id, null);
} catch {
// Race: another run grabbed the lock between our check and create.
await deps.prisma.agentRun.update({
where: { id: run.id },
data: { status: "FAILED", error: "lock race", finishedAt: new Date() },
});
await sendText(rt, chatId, "项目正在处理中,请稍候。");
return;
}
await sendText(rt, chatId, `已开始处理(model: ${model})。`);
const projectCtx: ProjectContext = {
projectId,
boundChatId: chatId, // ADR-0003: Feishu reads stay in this chat
workspaceDir: project.workspaceDir,
};
// Run the agent — fire-and-forget from the ws handler's perspective.
runAgent(deps.provider, deps.tools, { prompt, model, project: projectCtx })
.then(async (result) => {
await deps.prisma.agentRun.update({
where: { id: run.id },
data: {
status:
result.status === "completed"
? "COMPLETED"
: result.status === "length"
? "TIMED_OUT"
: "FAILED",
inputTokens: result.usage.inputTokens,
outputTokens: result.usage.outputTokens,
error: result.error ?? null,
finishedAt: new Date(),
},
});
await sendText(rt, chatId, statusReply(result.status));
})
.catch(async (e) => {
await deps.prisma.agentRun.update({
where: { id: run.id },
data: { status: "FAILED", error: String(e), finishedAt: new Date() },
});
await sendText(rt, chatId, "处理出错,已释放项目锁。");
})
.finally(async () => {
await releaseLock(deps.prisma, run.id);
});
};
}
function statusReply(status: "completed" | "length" | "failed"): string {
switch (status) {
case "completed":
return "处理完成。";
case "length":
return "处理达到步数上限,已停止。";
case "failed":
return "处理出错。";
}
}
/**
* Extract the prompt text from a text message, stripping @mentions.
* Returns null if the message is empty after stripping or doesn't mention the bot.
* Lark text content is JSON: `{"text":"@_user_1 do something"}`.
*/
function extractPrompt(msg: MessageReceiveEvent["message"]): string | null {
// Check for bot mention — mentions[].id.open_id should include the bot.
// The dispatcher already routes only message events; we further require
// that the bot is among the mentions (if no mentions, it's not @bot).
const mentions = msg.mentions;
if (mentions === undefined || mentions.length === 0) return null;
try {
const parsed = JSON.parse(msg.content) as { text?: string };
const text = parsed.text;
if (typeof text !== "string") return null;
// Strip @mentions (@_user_1 etc.) and trim; remainder is the prompt.
const stripped = text.replace(/@_\w+\s*/g, "").trim();
return stripped === "" ? null : stripped;
} catch {
return null;
}
}
+79
View File
@@ -0,0 +1,79 @@
/**
* ProjectAgentLock — ADR-0002 invariant enforcement at the DB layer.
*
* `projectId @id` gives at-most-one lock per project; `runId @unique` gives a
* run holds at most one lock. Those are DB-level. The spec's `WellFormed`
* invariant ("holder is a non-terminal run") is app-level: it must hold on
* every read of the lock table. `assertLockHolderActive` enforces it on the
* read path — if a lock exists but its run is terminal, that's a bug (the run
* should have released it), surfaced as an error rather than silently trusted.
*/
import type { PrismaClient } from "@prisma/client";
import type { AgentRunStatus } from "@prisma/client";
/// spec RunState.Terminal: completed/failed/timedOut/canceled.
/// active/waitingForUser are NOT terminal — the lock must still be held.
const TERMINAL_STATUSES: ReadonlySet<AgentRunStatus> = new Set([
"COMPLETED",
"FAILED",
"TIMED_OUT",
"CANCELED",
]);
export function isTerminal(status: AgentRunStatus): boolean {
return TERMINAL_STATUSES.has(status);
}
/**
* Acquire the project lock for a run. Throws if the project is already locked
* by another run (ADR-0002 exclusivity). The unique runId constraint means a
* run that already holds a lock elsewhere will also fail — that's correct, a
* run scopes to one project.
*/
export async function acquireLock(
prisma: PrismaClient,
projectId: string,
runId: string,
holderUserId: string | null,
): Promise<void> {
await prisma.projectAgentLock.create({
data: { projectId, runId, holderUserId },
});
}
/**
* Release the lock owned by a run. Idempotent — a no-op if no lock exists.
* Called on run completion/failure/timeout/cancellation (ADR-0002).
*/
export async function releaseLock(prisma: PrismaClient, runId: string): Promise<void> {
await prisma.projectAgentLock.deleteMany({ where: { runId } });
}
/**
* The spec's `LockTable.WellFormed` on a single read: if `projectId` has a
* lock owned by `runId`, that run must be non-terminal. Throws if a terminal
* run still holds a lock — that's a release-path bug, not a recoverable state.
*
* Returns the lock's runId if a healthy lock exists, or null if unlocked.
*/
export async function currentLockRunId(
prisma: PrismaClient,
projectId: string,
): Promise<string | null> {
const lock = await prisma.projectAgentLock.findUnique({
where: { projectId },
select: { runId: true },
});
if (lock === null) return null;
const run = await prisma.agentRun.findUnique({
where: { id: lock.runId },
select: { status: true },
});
if (run !== null && isTerminal(run.status)) {
throw new Error(
`lock invariant violated: project ${projectId} locked by terminal run ${lock.runId} (status ${run.status}) — release path bug (ADR-0002 WellFormed)`,
);
}
return lock.runId;
}
+58 -33
View File
@@ -1,31 +1,54 @@
/**
* Hub entry — skeleton.
* Hub entry — Fastify server + Feishu WebSocket listener.
*
* Wires the provider-agnostic agent seam end-to-end with a no-op Feishu read, so
* the invariant path is exercised without live credentials. Real wiring (Fastify
* server, Feishu OAuth/event callbacks, Prisma, admin-web) lands incrementally;
* this file proves the seam compiles and the pieces connect.
* Wires the full trigger path: lark ws receives `im.message.receive_v1` →
* trigger handler resolves chat→project (ADR-0001), creates AgentRun + acquires
* the project lock (ADR-0002), runs the provider-agnostic agent loop
* (ADR-0017), and releases the lock on finish. Feishu context reads inside the
* loop are ADR-0003-authorized (chat must match binding).
*
* Env (see .env.example):
* DATABASE_URL, OPENROUTER_API_KEY, FEISHU_APP_ID/SECRET/BOT_OPEN_ID, PORT
*/
import "dotenv/config";
import Fastify from "fastify";
import { prisma } from "./db.js";
import { OpenRouterProvider } from "./agent/openrouter-provider.js";
import { InMemoryModelRegistry } from "./agent/models.js";
import { ToolRegistry, feishuContextTool } from "./agent/tools.js";
import { runAgent } from "./agent/runner.js";
import { startFeishuListener } from "./feishu/client.js";
import { makeTriggerHandler } from "./feishu/trigger.js";
function main(): void {
const apiKey = process.env["OPENROUTER_API_KEY"];
if (apiKey === undefined || apiKey === "") {
console.error("[hub] OPENROUTER_API_KEY not set; nothing to run. Exiting.");
function requireEnv(name: string): string {
const v = process.env[name];
if (v === undefined || v === "") {
console.error(`[hub] missing required env: ${name}`);
process.exit(1);
}
return v;
}
const provider = new OpenRouterProvider({ apiKey });
async function main(): Promise<void> {
const databaseUrl = requireEnv("DATABASE_URL");
void databaseUrl; // prisma reads DATABASE_URL from env directly
const openrouterKey = requireEnv("OPENROUTER_API_KEY");
const feishuAppId = requireEnv("FEISHU_APP_ID");
const feishuAppSecret = requireEnv("FEISHU_APP_SECRET");
const feishuBotOpenId = requireEnv("FEISHU_BOT_OPEN_ID");
const port = Number(process.env["PORT"] ?? "8788");
const app = Fastify({ logger: true });
app.get("/api/healthz", async () => ({ ok: true, ts: Date.now() }));
// --- Agent seam wiring ---
const provider = new OpenRouterProvider({ apiKey: openrouterKey });
const models = new InMemoryModelRegistry(
[
{ id: "z-ai/glm-4.6", label: "GLM-4.6", toolCapable: true },
{ id: "anthropic/claude-3.5-sonnet", label: "Claude 3.5 Sonnet", toolCapable: true },
],
[
// Roles are data here — admin/teachers add presets, not code changes.
{ id: "draft", label: "草稿", defaultModel: "z-ai/glm-4.6" },
{ id: "review", label: "审校", defaultModel: "anthropic/claude-3.5-sonnet" },
],
@@ -33,30 +56,32 @@ function main(): void {
const tools = new ToolRegistry();
tools.register(
feishuContextTool(async (args) => {
// Skeleton: no live Feishu calls. Real impl wraps @larksuiteoapi/node-sdk.
return JSON.stringify({ stub: true, anchor: args.anchor, id: args.id });
// ADR-0003: the handler enforces chat==boundChat before any Feishu API call.
// Real impl wraps @larksuiteoapi/node-sdk's im.message.list/read APIs.
feishuContextTool(async (args, ctx) => {
return JSON.stringify({ stub: true, anchor: args.anchor, id: args.id, projectId: ctx.projectId });
}),
);
const model = models.resolve(undefined, "draft");
runAgent(provider, tools, {
prompt: "ping",
model,
project: {
projectId: "proj-skeleton",
boundChatId: "chat-skeleton",
workspaceDir: "/tmp/cph-skeleton",
},
})
.then((r) => {
console.log("[hub] run finished", { status: r.status, model, usage: r.usage });
if (r.error !== undefined) console.error("[hub] error:", r.error);
})
.catch((e) => {
console.error("[hub] run threw:", e);
// --- Feishu listener ---
const trigger = makeTriggerHandler({ prisma, provider, tools, models, logger: app.log });
const feishuCtx = startFeishuListener(
{ appId: feishuAppId, appSecret: feishuAppSecret, botOpenId: feishuBotOpenId },
app.log,
trigger,
);
void feishuCtx;
app.listen({ port, host: "0.0.0.0" }, (err, address) => {
if (err !== null) {
app.log.error({ err }, "listen failed");
process.exit(1);
});
}
app.log.info({ address }, "hub listening");
});
}
main();
main().catch((e) => {
console.error("[hub] fatal:", e);
process.exit(1);
});