diff --git a/hub/src/feishu/trigger.ts b/hub/src/feishu/trigger.ts index 657ca95..da5f9cd 100644 --- a/hub/src/feishu/trigger.ts +++ b/hub/src/feishu/trigger.ts @@ -37,9 +37,12 @@ import { readFeishuContext } from "./read.js"; import { MessageBatcher, messageBatchKey, type MessageBatcherOptions } from "./messageBatcher.js"; import { ApprovalManager } from "./approval.js"; import { SenderNameCache } from "./senderCache.js"; +import { TriggerQueue, triggerQueue as defaultTriggerQueue, type QueuedTrigger } from "./triggerQueue.js"; export { ApprovalManager } from "./approval.js"; export type { ApprovalResult, PendingApproval } from "./approval.js"; +export { TriggerQueue, triggerQueue } from "./triggerQueue.js"; +export type { QueuedTrigger, TriggerQueueOptions } from "./triggerQueue.js"; export const senderNameCache = new SenderNameCache(); interface TriggerDeps { @@ -49,6 +52,7 @@ interface TriggerDeps { readonly runAgent?: (req: RunRequest) => Promise; readonly authorizer?: PermissionAuthorizer | undefined; readonly messageBatcherOptions?: MessageBatcherOptions | undefined; + readonly triggerQueue?: TriggerQueue | undefined; } interface TriggerActor { @@ -65,6 +69,8 @@ interface TriggerRunContext { readonly actor: TriggerActor; } +type StartAgentRunOutcome = "started" | "queued" | "rejected" | "locked" | "skipped"; + export interface TriggerHandler { (event: MessageReceiveEvent, rt: FeishuRuntime): Promise; readonly onCardAction: (event: CardActionEvent, rt: FeishuRuntime) => Promise; @@ -80,6 +86,7 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { const authorizer = deps.authorizer ?? createPermissionAuthorizer(deps.prisma); const runAgent = deps.runAgent ?? defaultRunAgent; const approvalManager = new ApprovalManager(); + const triggerQueue = deps.triggerQueue ?? defaultTriggerQueue; const batchContexts = new Map(); const messageBatcher = new MessageBatcher(async (mergedText, key) => { if (key === undefined) { @@ -100,15 +107,19 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { } }, deps.messageBatcherOptions); - async function startAgentRun(context: TriggerRunContext, cleanPrompt: string): Promise { + async function startAgentRun( + context: TriggerRunContext, + cleanPrompt: string, + options: { readonly queueIfLocked?: boolean } = {}, + ): Promise { const { msg, rt, chatId, projectId, senderOpenId, actor } = context; + const queueIfLocked = options.queueIfLocked ?? true; - // ADR-0002: if locked by another run, reply busy. + // ADR-0002: if locked by another run, enqueue the trigger for this project. const existing = await currentLockRunId(deps.prisma, projectId); if (existing !== null) { - await reactToMessage(rt, msg.message_id, "OnIt"); - await sendText(rt, chatId, "项目正在处理中,请稍候。"); - return; + if (!queueIfLocked) return "locked"; + return enqueueLockedTrigger(context, cleanPrompt); } const project = await deps.prisma.project.findUnique({ @@ -117,7 +128,7 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { }); if (project === null) { deps.logger.error({ projectId }, "feishu trigger: project missing"); - return; + return "skipped"; } // ADR-0017: role-as-data. Parse a leading `/` command; unknown @@ -148,7 +159,7 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { }, }); await sendText(rt, chatId, `无权限使用角色 ${roleId}。`); - return; + return "skipped"; } const role = models.role(roleId); const model = models.resolve(undefined, roleId); @@ -234,9 +245,8 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { action: "run.lock_race", metadata: { sender: await senderAuditMetadata(rt, senderOpenId) }, }); - await reactToMessage(rt, msg.message_id, "OnIt"); - await sendText(rt, chatId, "项目正在处理中,请稍候。"); - return; + if (!queueIfLocked) return "locked"; + return enqueueLockedTrigger(context, cleanPrompt); } const projectCtx: ProjectContext = { @@ -344,7 +354,85 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { } catch (e) { deps.logger.warn({ runId: run.id, err: e instanceof Error ? e.message : String(e) }, "trigger: could not release lock"); } + try { + await drainTriggerQueue(projectId, rt); + } catch (e) { + deps.logger.warn({ projectId, err: e instanceof Error ? e.message : String(e) }, "trigger: could not drain trigger queue"); + } }); + return "started"; + } + + async function enqueueLockedTrigger(context: TriggerRunContext, cleanPrompt: string): Promise { + const position = triggerQueue.enqueue(context.projectId, { + chatId: context.chatId, + prompt: cleanPrompt, + msg: context.msg, + senderOpenId: context.senderOpenId, + actor: context.actor, + }); + if (position === 0) { + await reactToMessage(context.rt, context.msg.message_id, "OnIt"); + await sendText(context.rt, context.chatId, "队列已满,请稍后再试"); + return "rejected"; + } + + await sendText(context.rt, context.chatId, `已加入队列(第${position}位),当前处理完成后将自动开始`); + return "queued"; + } + + async function drainTriggerQueue(projectId: string, rt: FeishuRuntime): Promise { + while (triggerQueue.hasPending(projectId)) { + const existing = await currentLockRunId(deps.prisma, projectId); + if (existing !== null) return; + + const queued = triggerQueue.peek(projectId); + if (queued === null) return; + if (triggerQueue.isExpired(queued)) { + triggerQueue.dequeue(projectId); + deps.logger.info( + { projectId, messageId: queued.msg.message_id, enqueuedAt: queued.enqueuedAt }, + "feishu trigger: dropped expired queued trigger", + ); + continue; + } + + const next = triggerQueue.dequeue(projectId); + if (next === null) return; + const outcome = await startAgentRun(contextFromQueuedTrigger(next, rt), next.prompt, { queueIfLocked: false }); + if (outcome === "started") return; + if (outcome === "locked") { + requeueTrigger(next); + return; + } + } + } + + function contextFromQueuedTrigger(trigger: QueuedTrigger, rt: FeishuRuntime): TriggerRunContext { + return { + msg: trigger.msg, + rt, + chatId: trigger.chatId, + projectId: trigger.projectId, + senderOpenId: trigger.senderOpenId, + actor: trigger.actor, + }; + } + + function requeueTrigger(trigger: QueuedTrigger): void { + const position = triggerQueue.enqueue(trigger.projectId, { + chatId: trigger.chatId, + prompt: trigger.prompt, + msg: trigger.msg, + senderOpenId: trigger.senderOpenId, + actor: trigger.actor, + }); + if (position === 0) { + deps.logger.warn( + { projectId: trigger.projectId, messageId: trigger.msg.message_id }, + "feishu trigger: could not requeue locked queued trigger", + ); + } } const onCardAction = async (event: CardActionEvent, rt: FeishuRuntime): Promise => { @@ -479,7 +567,7 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { // Slash commands: session management, not agent runs. These bypass the // batcher. Session commands bypass the lock because they don't create runs; - // role/unknown slash prompts still run immediately through the normal lock. + // role/unknown slash prompts still use the normal run path and queue if locked. if (cleanPrompt.startsWith("/")) { const cmd = cleanPrompt.split(/\s+/)[0]; switch (cmd) { @@ -513,6 +601,10 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { where: { projectId, archivedAt: null }, data: { archivedAt: new Date() }, }); + const cleared = triggerQueue.clear(projectId); + if (cleared > 0) { + deps.logger.info({ projectId, cleared }, "feishu trigger: cleared queued triggers on reset"); + } await sendText(rt, chatId, "已重置,下次 @bot 将从头开始。"); return; } @@ -526,8 +618,7 @@ export function makeTriggerHandler(deps: TriggerDeps): TriggerHandler { if (msg.message_type === "text") { const existing = await currentLockRunId(deps.prisma, projectId); if (existing !== null) { - await reactToMessage(rt, msg.message_id, "OnIt"); - await sendText(rt, chatId, "项目正在处理中,请稍候。"); + await enqueueLockedTrigger(runContext, cleanPrompt); return; } const key = messageBatchKey(chatId, senderOpenId); diff --git a/hub/src/feishu/triggerQueue.ts b/hub/src/feishu/triggerQueue.ts new file mode 100644 index 0000000..ae86f8c --- /dev/null +++ b/hub/src/feishu/triggerQueue.ts @@ -0,0 +1,100 @@ +import type { MessageReceiveEvent } from "./client.js"; + +export interface QueuedTrigger { + readonly projectId: string; + readonly chatId: string; + readonly prompt: string; + readonly msg: MessageReceiveEvent["message"]; + readonly senderOpenId: string; + readonly actor: { readonly feishuOpenId: string; readonly chatId: string }; + readonly enqueuedAt: number; +} + +export interface TriggerQueueOptions { + readonly maxQueueSize?: number; + readonly maxWaitMs?: number; +} + +const DEFAULT_MAX_QUEUE_SIZE = 5; +const DEFAULT_MAX_WAIT_MS = 300_000; + +export class TriggerQueue { + readonly maxQueueSize: number; + readonly maxWaitMs: number; + + private readonly queues = new Map(); + + constructor(options: TriggerQueueOptions = {}) { + this.maxQueueSize = Math.max(0, Math.trunc(options.maxQueueSize ?? DEFAULT_MAX_QUEUE_SIZE)); + this.maxWaitMs = Math.max(0, Math.trunc(options.maxWaitMs ?? DEFAULT_MAX_WAIT_MS)); + } + + enqueue(projectId: string, trigger: Omit): number { + const queue = this.queues.get(projectId) ?? []; + if (queue.length >= this.maxQueueSize) return 0; + + queue.push({ + ...trigger, + projectId, + enqueuedAt: Date.now(), + }); + if (!this.queues.has(projectId)) { + this.queues.set(projectId, queue); + } + return queue.length; + } + + dequeue(projectId: string): QueuedTrigger | null { + const queue = this.queues.get(projectId); + if (queue === undefined || queue.length === 0) return null; + + const next = queue.shift() ?? null; + if (queue.length === 0) { + this.queues.delete(projectId); + } + return next; + } + + peek(projectId: string): QueuedTrigger | null { + return this.queues.get(projectId)?.[0] ?? null; + } + + length(projectId: string): number { + return this.queues.get(projectId)?.length ?? 0; + } + + hasPending(projectId: string): boolean { + return this.length(projectId) > 0; + } + + purgeExpired(): number { + const now = Date.now(); + let removed = 0; + for (const [projectId, queue] of this.queues) { + const fresh = queue.filter((trigger) => !this.isExpired(trigger, now)); + removed += queue.length - fresh.length; + if (fresh.length === 0) { + this.queues.delete(projectId); + } else if (fresh.length !== queue.length) { + this.queues.set(projectId, fresh); + } + } + return removed; + } + + clear(projectId: string): number { + const count = this.length(projectId); + this.queues.delete(projectId); + return count; + } + + clearAll(): void { + this.queues.clear(); + } + + isExpired(trigger: QueuedTrigger, now = Date.now()): boolean { + return trigger.enqueuedAt + this.maxWaitMs < now; + } +} + +export const triggerQueue = new TriggerQueue(); diff --git a/hub/src/server.ts b/hub/src/server.ts index e6543bc..4612b37 100644 --- a/hub/src/server.ts +++ b/hub/src/server.ts @@ -16,6 +16,7 @@ import { prisma } from "./db.js"; import { createEnvRuntimeSettings } from "./settings/runtime.js"; import { createLarkClient, startFeishuListenerWithClient } from "./feishu/client.js"; import { makeTriggerHandler } from "./feishu/trigger.js"; +import { triggerQueue } from "./feishu/triggerQueue.js"; function requireEnv(name: string): string { const v = process.env[name]; @@ -61,6 +62,13 @@ async function main(): Promise { if (feishuListenerEnabled) { const feishuConfig = { appId: feishuAppId, appSecret: feishuAppSecret, botOpenId: feishuBotOpenId }; const larkClient = createLarkClient(feishuConfig); + const triggerQueuePurgeTimer = setInterval(() => { + const removed = triggerQueue.purgeExpired(); + if (removed > 0) { + app.log.info({ removed }, "feishu trigger queue: purged expired triggers"); + } + }, 60_000); + triggerQueuePurgeTimer.unref(); const trigger = makeTriggerHandler({ prisma, settings: runtimeSettings, logger: app.log }); startFeishuListenerWithClient(feishuConfig, larkClient, app.log, trigger, trigger.onCardAction); } else { diff --git a/hub/test/integration/trigger.test.ts b/hub/test/integration/trigger.test.ts index e616c81..d61260a 100644 --- a/hub/test/integration/trigger.test.ts +++ b/hub/test/integration/trigger.test.ts @@ -2,6 +2,7 @@ import { describe, it, expect, beforeEach, afterAll, vi } from "vitest"; import { prisma, resetDb, mockFeishuRuntime, seedProject, silentLogger } from "./helpers.js"; import { InMemoryModelRegistry } from "../../src/agent/models.js"; import { makeTriggerHandler, extractPrompt } from "../../src/feishu/trigger.js"; +import { TriggerQueue } from "../../src/feishu/triggerQueue.js"; import type { MessageReceiveEvent } from "../../src/feishu/client.js"; import type { RunRequest, RunResult } from "../../src/agent/runner.js"; import type { RuntimeSettings } from "../../src/settings/runtime.js"; @@ -179,7 +180,7 @@ describe("trigger full lifecycle (integration)", () => { expect(runs).toHaveLength(0); }); - it("replies busy when project is already locked (ADR-0002)", async () => { + it("queues a text trigger when project is already locked (ADR-0002)", async () => { await seedProject("proj-3", "chat-3"); // Manually create a lock by inserting a run + lock. const existingRun = await prisma.agentRun.create({ @@ -192,13 +193,59 @@ describe("trigger full lifecycle (integration)", () => { const trigger = makeTriggerHandler({ prisma, settings, logger: silentLogger, runAgent, messageBatcherOptions: { maxMessages: 1 } }); await trigger(makeEvent("chat-3", "@_user_1 写教案"), rt); - expect(rt.sentTexts).toContain("项目正在处理中,请稍候。"); + expect(rt.sentTexts).toContain("已加入队列(第1位),当前处理完成后将自动开始"); + expect(rt.reactions.some((reaction) => reaction.emoji === "OnIt")).toBe(false); expect(runAgentCalls).toHaveLength(0); // No new run created. const runs = await prisma.agentRun.findMany(); expect(runs).toHaveLength(1); }); + it("starts the next queued text trigger when the current run finishes", async () => { + await seedProject("proj-3b", "chat-3b"); + const firstRun = deferred(); + const secondRun = deferred(); + const pendingRuns = [firstRun, secondRun]; + const queuedRunAgent: TestRunner = async (req) => { + runAgentCalls.push(req); + req.onStream?.({ type: "text-delta", text: `mock response ${runAgentCalls.length}` }); + const pendingRun = pendingRuns.shift(); + if (pendingRun === undefined) { + throw new Error("unexpected extra run"); + } + return pendingRun.promise; + }; + const trigger = makeTriggerHandler({ + prisma, + settings, + logger: silentLogger, + runAgent: queuedRunAgent, + messageBatcherOptions: { maxMessages: 1 }, + }); + + await trigger(makeEvent("chat-3b", "@_user_1 第一个请求"), rt); + await vi.waitFor(() => { + expect(runAgentCalls).toHaveLength(1); + }); + + await trigger(makeEvent("chat-3b", "@_user_1 第二个请求"), rt); + expect(rt.sentTexts).toContain("已加入队列(第1位),当前处理完成后将自动开始"); + expect(runAgentCalls).toHaveLength(1); + + firstRun.resolve(completedRunResult("first done", "sdk-session-first")); + await vi.waitFor(() => { + expect(runAgentCalls).toHaveLength(2); + }); + expect(runAgentCalls[1]?.prompt).toBe("第二个请求"); + + secondRun.resolve(completedRunResult("second done", "sdk-session-second")); + await vi.waitFor(async () => { + const runs = await prisma.agentRun.findMany({ orderBy: { createdAt: "asc" } }); + expect(runs).toHaveLength(2); + expect(runs.every((run) => run.status === "COMPLETED")).toBe(true); + }); + }); + it("ignores messages from unbound chats (ADR-0001)", async () => { await seedProject("proj-4", "chat-4"); const trigger = makeTriggerHandler({ prisma, settings, logger: silentLogger, runAgent, messageBatcherOptions: { maxMessages: 1 } }); @@ -278,7 +325,15 @@ describe("trigger full lifecycle (integration)", () => { it("/reset archives current session", async () => { await seedProject("proj-8", "chat-8"); - const trigger = makeTriggerHandler({ prisma, settings, logger: silentLogger, runAgent, messageBatcherOptions: { maxMessages: 1 } }); + const queue = new TriggerQueue(); + const trigger = makeTriggerHandler({ + prisma, + settings, + logger: silentLogger, + runAgent, + messageBatcherOptions: { maxMessages: 1 }, + triggerQueue: queue, + }); await trigger(makeEvent("chat-8", "@_user_1 写教案"), rt); await vi.waitFor(async () => { @@ -287,8 +342,19 @@ describe("trigger full lifecycle (integration)", () => { expect(runs[0]?.status).toBe("COMPLETED"); }); + const queuedEvent = makeEvent("chat-8", "@_user_1 后续需求"); + queue.enqueue("proj-8", { + chatId: "chat-8", + prompt: extractPrompt(queuedEvent.message) ?? "后续需求", + msg: queuedEvent.message, + senderOpenId: "ou_test_user", + actor: { feishuOpenId: "ou_test_user", chatId: "chat-8" }, + }); + expect(queue.length("proj-8")).toBe(1); + await trigger(makeEvent("chat-8", "@_user_1 /reset"), rt); expect(rt.sentTexts).toContain("已重置,下次 @bot 将从头开始。"); + expect(queue.length("proj-8")).toBe(0); const sessions = await prisma.agentSession.findMany(); expect(sessions).toHaveLength(1); @@ -490,6 +556,32 @@ describe("trigger full lifecycle (integration)", () => { }); }); +interface Deferred { + readonly promise: Promise; + readonly resolve: (value: T) => void; + readonly reject: (reason?: unknown) => void; +} + +function deferred(): Deferred { + let resolve!: (value: T) => void; + let reject!: (reason?: unknown) => void; + const promise = new Promise((res, rej) => { + resolve = res; + reject = rej; + }); + return { promise, resolve, reject }; +} + +function completedRunResult(text: string, sdkSessionId: string): RunResult { + return { + status: "completed", + text, + usage: { inputTokens: 10, outputTokens: 5 }, + numTurns: 1, + sdkSessionId, + }; +} + afterAll(async () => { await prisma.$disconnect(); }); diff --git a/hub/test/unit/trigger-queue.test.ts b/hub/test/unit/trigger-queue.test.ts new file mode 100644 index 0000000..e54beda --- /dev/null +++ b/hub/test/unit/trigger-queue.test.ts @@ -0,0 +1,135 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { TriggerQueue, type QueuedTrigger } from "../../src/feishu/triggerQueue.js"; +import type { MessageReceiveEvent } from "../../src/feishu/client.js"; + +describe("TriggerQueue", () => { + afterEach(() => { + vi.useRealTimers(); + }); + + it("dequeues triggers in FIFO order", () => { + const queue = new TriggerQueue(); + + queue.enqueue("project-1", makeTrigger("first")); + queue.enqueue("project-1", makeTrigger("second")); + + expect(queue.dequeue("project-1")?.prompt).toBe("first"); + expect(queue.dequeue("project-1")?.prompt).toBe("second"); + }); + + it("returns the 1-based queue position from enqueue", () => { + const queue = new TriggerQueue(); + + expect(queue.enqueue("project-1", makeTrigger("first"))).toBe(1); + expect(queue.enqueue("project-1", makeTrigger("second"))).toBe(2); + expect(queue.enqueue("project-1", makeTrigger("third"))).toBe(3); + }); + + it("returns 0 when the project queue is full", () => { + const queue = new TriggerQueue({ maxQueueSize: 2 }); + + expect(queue.enqueue("project-1", makeTrigger("first"))).toBe(1); + expect(queue.enqueue("project-1", makeTrigger("second"))).toBe(2); + expect(queue.enqueue("project-1", makeTrigger("third"))).toBe(0); + expect(queue.length("project-1")).toBe(2); + }); + + it("returns null when dequeueing an empty queue", () => { + const queue = new TriggerQueue(); + + expect(queue.dequeue("project-1")).toBeNull(); + }); + + it("peeks without removing the trigger", () => { + const queue = new TriggerQueue(); + + queue.enqueue("project-1", makeTrigger("first")); + + expect(queue.peek("project-1")?.prompt).toBe("first"); + expect(queue.length("project-1")).toBe(1); + expect(queue.dequeue("project-1")?.prompt).toBe("first"); + }); + + it("reports length and pending state per project", () => { + const queue = new TriggerQueue(); + + expect(queue.length("project-1")).toBe(0); + expect(queue.hasPending("project-1")).toBe(false); + + queue.enqueue("project-1", makeTrigger("first")); + + expect(queue.length("project-1")).toBe(1); + expect(queue.hasPending("project-1")).toBe(true); + }); + + it("purges expired triggers", () => { + vi.useFakeTimers(); + vi.setSystemTime(1_000); + const queue = new TriggerQueue({ maxWaitMs: 100 }); + + queue.enqueue("project-1", makeTrigger("old")); + vi.setSystemTime(1_050); + queue.enqueue("project-1", makeTrigger("fresh")); + vi.setSystemTime(1_101); + + expect(queue.purgeExpired()).toBe(1); + expect(queue.length("project-1")).toBe(1); + expect(queue.dequeue("project-1")?.prompt).toBe("fresh"); + }); + + it("clears all triggers for a project", () => { + const queue = new TriggerQueue(); + + queue.enqueue("project-1", makeTrigger("first")); + queue.enqueue("project-1", makeTrigger("second")); + queue.enqueue("project-2", makeTrigger("other")); + + expect(queue.clear("project-1")).toBe(2); + expect(queue.length("project-1")).toBe(0); + expect(queue.length("project-2")).toBe(1); + }); + + it("clears all project queues", () => { + const queue = new TriggerQueue(); + + queue.enqueue("project-1", makeTrigger("first")); + queue.enqueue("project-2", makeTrigger("other")); + queue.clearAll(); + + expect(queue.hasPending("project-1")).toBe(false); + expect(queue.hasPending("project-2")).toBe(false); + }); + + it("keeps multiple project queues independent", () => { + const queue = new TriggerQueue(); + + expect(queue.enqueue("project-1", makeTrigger("p1-first"))).toBe(1); + expect(queue.enqueue("project-2", makeTrigger("p2-first"))).toBe(1); + expect(queue.enqueue("project-1", makeTrigger("p1-second"))).toBe(2); + + expect(queue.dequeue("project-2")?.prompt).toBe("p2-first"); + expect(queue.dequeue("project-1")?.prompt).toBe("p1-first"); + expect(queue.dequeue("project-1")?.prompt).toBe("p1-second"); + }); +}); + +function makeTrigger(prompt: string): Omit { + return { + chatId: "chat-1", + prompt, + msg: makeMessage(prompt), + senderOpenId: "ou-user", + actor: { feishuOpenId: "ou-user", chatId: "chat-1" }, + }; +} + +function makeMessage(prompt: string): MessageReceiveEvent["message"] { + return { + message_id: `message-${prompt}`, + chat_id: "chat-1", + chat_type: "group", + message_type: "text", + content: JSON.stringify({ text: `@_user_1 ${prompt}` }), + mentions: [{ key: "@_user_1", id: { open_id: "ou-bot" }, name: "Bot" }], + }; +}