feat: add per-project message trigger wait queue

- TriggerQueue: FIFO per-project queue, max 5 items, 5min TTL expiry
- Replace reject-when-locked with queue: users get position feedback
- Run completion auto-starts next queued trigger (skips expired items)
- /reset clears project queue; 60s periodic purge interval
- Add unit tests (10) for queue lifecycle and edge cases
This commit is contained in:
2026-07-08 17:31:43 +08:00
parent 4479f88f4f
commit 36c0801e6f
5 changed files with 442 additions and 16 deletions
+104 -13
View File
@@ -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<RunResult>;
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<void>;
readonly onCardAction: (event: CardActionEvent, rt: FeishuRuntime) => Promise<void>;
@@ -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<string, TriggerRunContext>();
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<void> {
async function startAgentRun(
context: TriggerRunContext,
cleanPrompt: string,
options: { readonly queueIfLocked?: boolean } = {},
): Promise<StartAgentRunOutcome> {
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 `/<role>` 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<StartAgentRunOutcome> {
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<void> {
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<void> => {
@@ -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);
+100
View File
@@ -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<string, QueuedTrigger[]>();
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<QueuedTrigger, "projectId" | "enqueuedAt">): 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();
+8
View File
@@ -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<void> {
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 {
+95 -3
View File
@@ -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<RunResult>();
const secondRun = deferred<RunResult>();
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<T> {
readonly promise: Promise<T>;
readonly resolve: (value: T) => void;
readonly reject: (reason?: unknown) => void;
}
function deferred<T>(): Deferred<T> {
let resolve!: (value: T) => void;
let reject!: (reason?: unknown) => void;
const promise = new Promise<T>((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();
});
+135
View File
@@ -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<QueuedTrigger, "projectId" | "enqueuedAt"> {
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" }],
};
}