feat: add assistant proposal worker runtime

This commit is contained in:
zenord
2026-08-18 18:48:47 +08:00
parent 947f19e3cd
commit b4121c9fd6
32 changed files with 2827 additions and 986 deletions
+87 -196
View File
@@ -3,35 +3,13 @@ import type { GatewayPolicy } from "../config.js";
import type { PlatformAdapter } from "./adapter.js";
import { CommandRouter, type ParsedCommand } from "./command-router.js";
import { chatKeyFor } from "./durable-session-store.js";
import type { IncomingMessage } from "./types.js";
import type { IncomingMessage, MessageTarget } from "./types.js";
export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string }
export interface GatewayTimer { cancel(): void }
export interface GatewayOptions {
now?: () => number;
schedule?: (callback: () => void, delayMs: number) => GatewayTimer;
}
interface ChatWorkState {
running: boolean;
queued: number;
startedAt?: number;
}
interface AsyncInbound {
phase: "queued" | "running" | "finished";
receivedAt: number;
initialNoticeSent: boolean;
queuedNoticeSent: boolean;
timer?: GatewayTimer;
stream: ReplyStream;
}
interface LockedResult {
result: GatewayResult;
delivery: Promise<void>;
interface ChatTarget {
target: MessageTarget;
adapter: PlatformAdapter;
}
class ReplyStream {
@@ -60,22 +38,14 @@ class ReplyStream {
}
export class Gateway {
private readonly locks = new Map<string, Promise<void>>();
private readonly chatWork = new Map<string, ChatWorkState>();
private readonly asyncInbounds = new Set<AsyncInbound>();
private readonly now: () => number;
private readonly schedule: (callback: () => void, delayMs: number) => GatewayTimer;
private readonly chatTargets = new Map<string, ChatTarget>();
private closed = false;
readonly commandRouter = new CommandRouter();
constructor(
private readonly policy: GatewayPolicy,
private readonly runtime: ConversationRuntime,
options: GatewayOptions = {}
) {
this.now = options.now || Date.now;
this.schedule = options.schedule || defaultSchedule;
}
private readonly runtime: ConversationRuntime
) {}
async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise<GatewayResult> {
const policyError = this.checkPolicy(message);
@@ -84,82 +54,74 @@ export class Gateway {
return { ok: true, ignored: true, error: policyError };
}
const command = this.commandRouter.parse(message.text);
if (command?.kind === "help" || command?.kind === "status" || command?.kind === "cancel") {
return this.executeAndReply(command, message, adapter, options);
}
this.chatTargets.set(chatKeyFor(message.platform, message.chatId), {
target: {
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
raw: message.raw
},
adapter
});
const chatKey = chatKeyFor(message.platform, message.chatId);
if (command || options.synchronous) {
const result = await this.withChatLock(chatKey, () => this.executeSynchronous(command, message));
const command = this.commandRouter.parse(message.text);
if (command) {
const result = await this.executeCommandResult(command, message);
if (!options.synchronous) await new ReplyStream(message, adapter).enqueue(result.reply || "");
return result;
}
const inbound: AsyncInbound = {
phase: "queued",
receivedAt: this.now(),
initialNoticeSent: false,
queuedNoticeSent: false,
stream: new ReplyStream(message, adapter)
};
this.asyncInbounds.add(inbound);
this.armTimer(inbound, inbound.receivedAt + 15_000);
if (options.synchronous) return this.promptResult(message);
const locked = await this.withChatLock(chatKey, async () => {
inbound.phase = "running";
if (inbound.queuedNoticeSent) void inbound.stream.enqueue("轮到这件事了,我现在开始处理。");
try {
const reply = (await this.runtime.prompt({
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
text: message.text,
messageId: message.messageId
})).text;
this.finishInbound(inbound);
return { result: { ok: true, reply }, delivery: inbound.stream.enqueue(reply) };
} catch (error) {
const errorText = error instanceof Error ? error.message : String(error);
const reply = `Agent error: ${errorText}`;
this.finishInbound(inbound);
return { result: { ok: false, error: errorText, reply }, delivery: inbound.stream.enqueue(reply) };
}
});
await locked.delivery;
return locked.result;
try {
const reply = (await this.runtime.prompt({
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
text: message.text,
messageId: message.messageId
})).text;
await new ReplyStream(message, adapter).enqueue(reply);
return { ok: true, reply };
} catch (error) {
const errorText = error instanceof Error ? error.message : String(error);
const reply = `Agent error: ${errorText}`;
await new ReplyStream(message, adapter).enqueue(reply);
return { ok: false, error: errorText, reply };
}
}
stats(): ReturnType<ConversationRuntime["stats"]> & { lockedChats: number } {
return { ...this.runtime.stats(), lockedChats: this.locks.size };
async sendEvent(chatKey: string, text: string): Promise<void> {
if (this.closed) return;
const entry = this.chatTargets.get(chatKey);
if (!entry) return;
try {
await entry.adapter.sendMessage({ target: entry.target, text });
} catch {
console.error(`Gateway event send failed (platform=${entry.target.platform})`);
}
}
stats(): ReturnType<ConversationRuntime["stats"]> {
return this.runtime.stats();
}
shutdown(): void {
if (this.closed) return;
this.closed = true;
for (const inbound of this.asyncInbounds) {
inbound.phase = "finished";
inbound.timer?.cancel();
inbound.timer = undefined;
}
this.asyncInbounds.clear();
}
private async executeAndReply(
command: ParsedCommand,
message: IncomingMessage,
adapter: PlatformAdapter,
options: { synchronous?: boolean }
): Promise<GatewayResult> {
const result = await this.executeSynchronous(command, message);
if (options.synchronous) return result;
await new ReplyStream(message, adapter).enqueue(result.reply || "");
return result;
}
private async executeSynchronous(command: ParsedCommand | undefined, message: IncomingMessage): Promise<GatewayResult> {
private async executeCommandResult(command: ParsedCommand, message: IncomingMessage): Promise<GatewayResult> {
try {
const reply = command ? await this.executeCommand(command, message) : (await this.runtime.prompt({
return { ok: true, reply: await this.executeCommand(command, message) };
} catch (error) {
const errorText = error instanceof Error ? error.message : String(error);
return { ok: false, error: errorText, reply: `Agent error: ${errorText}` };
}
}
private async promptResult(message: IncomingMessage): Promise<GatewayResult> {
try {
const reply = (await this.runtime.prompt({
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
@@ -175,60 +137,30 @@ export class Gateway {
private async executeCommand(command: ParsedCommand, message: IncomingMessage): Promise<string> {
switch (command.kind) {
case "help": return ["Commands:", "/status", "/cancel", "/new", "/help"].join("\n");
case "help": return HELP_TEXT;
case "status": {
const chatKey = chatKeyFor(message.platform, message.chatId);
const work = this.chatWork.get(chatKey);
const status = {
...this.runtime.status(message.platform, message.chatId),
gatewayRunning: Boolean(work?.running),
queued: work?.queued || 0,
gatewayRunningSeconds: work?.startedAt !== undefined
? Math.max(0, Math.floor((this.now() - work.startedAt) / 1000)) : 0
};
const status = this.runtime.status(message.platform, message.chatId);
return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`;
}
case "new":
await this.runtime.reset(message.platform, message.chatId);
return "Started a new native ACP session for this chat.";
case "cancel": return this.cancelText(message);
case "retired-role": return "Role switching was removed in Config v3; this instance has one fixed Bot and ACP agent.";
case "list": {
if (!this.runtime.listProposals) return "当前运行时不支持列出提案。";
const proposals = this.runtime.listProposals(message.platform, message.chatId);
if (proposals.length === 0) return "当前没有提案。";
return proposals.map((proposal) => `${proposal.id} [${proposal.status}] ${proposal.title}`).join("\n");
}
case "confirm": {
if (!this.runtime.confirm) return "当前运行时不支持确认操作。";
return await this.runtime.confirm(message.platform, message.chatId) ? "已确认,继续执行。" : "当前没有待确认的提案。";
}
case "stop": {
if (!this.runtime.stop) return "当前运行时不支持停止操作。";
return await this.runtime.stop(message.platform, message.chatId) ? "已停止当前任务。" : "当前没有正在执行的任务。";
}
case "cancel":
return await this.runtime.cancel(message.platform, message.chatId) ? "已取消最近的提案。" : "没有可取消的提案。";
}
}
private async cancelText(message: IncomingMessage): Promise<string> {
return await this.runtime.cancel(message.platform, message.chatId) ? "Cancellation requested." : "No active turn to cancel.";
}
private armTimer(inbound: AsyncInbound, deadline: number): void {
if (this.closed || inbound.phase === "finished") return;
inbound.timer?.cancel();
inbound.timer = this.schedule(() => {
inbound.timer = undefined;
if (this.closed || inbound.phase === "finished") return;
const now = this.now();
if (now < deadline) {
this.armTimer(inbound, deadline);
return;
}
const elapsed = now - inbound.receivedAt;
if (!inbound.initialNoticeSent) {
inbound.initialNoticeSent = true;
inbound.queuedNoticeSent = inbound.phase === "queued";
}
void inbound.stream.enqueue(reminderText(inbound.phase, elapsed));
this.armTimer(inbound, nextReminderDeadline(inbound.receivedAt, now));
}, Math.max(0, deadline - this.now()));
}
private finishInbound(inbound: AsyncInbound): void {
if (inbound.phase === "finished") return;
inbound.phase = "finished";
inbound.timer?.cancel();
inbound.timer = undefined;
this.asyncInbounds.delete(inbound);
}
private checkPolicy(message: IncomingMessage): string | undefined {
if (this.policy.allowedUsers.length > 0 && !this.policy.allowedUsers.includes(message.userId)) return "User not allowed";
if (this.policy.allowedChats.length > 0 && !this.policy.allowedChats.includes(message.chatId)) return "Chat not allowed";
@@ -236,55 +168,14 @@ export class Gateway {
return undefined;
}
private async withChatLock<T>(key: string, fn: () => Promise<T>): Promise<T> {
const previous = this.locks.get(key) || Promise.resolve();
const state = this.chatWork.get(key) || { running: false, queued: 0 };
state.queued++;
this.chatWork.set(key, state);
let release!: () => void;
const current = new Promise<void>((resolve) => { release = resolve; });
const lock = previous.then(() => current);
this.locks.set(key, lock);
await previous;
state.queued--;
state.running = true;
state.startedAt = this.now();
try {
return await fn();
} finally {
state.running = false;
state.startedAt = undefined;
release();
if (this.locks.get(key) === lock) {
this.locks.delete(key);
this.chatWork.delete(key);
}
}
}
}
function defaultSchedule(callback: () => void, delayMs: number): GatewayTimer {
const timer = setTimeout(callback, delayMs);
timer.unref();
return { cancel: () => clearTimeout(timer) };
}
function nextReminderDeadline(receivedAt: number, now: number): number {
const elapsed = now - receivedAt;
if (elapsed < 60_000) return receivedAt + 60_000;
if (elapsed < 180_000) return receivedAt + 180_000;
if (elapsed < 480_000) return receivedAt + 480_000;
return receivedAt + 480_000 + (Math.floor((elapsed - 480_000) / 600_000) + 1) * 600_000;
}
function reminderText(phase: AsyncInbound["phase"], elapsed: number): string {
if (phase === "queued") {
return elapsed < 60_000
? "前面还有一件事没处理完,这条我记着,轮到后马上处理。"
: "前面的事情还没处理完,这条还在等。我没有漏掉,轮到后会接着做。";
}
if (elapsed < 60_000) return "收到,我还在处理,稍等我一下。";
if (elapsed < 180_000) return "还在弄,暂时还没出结果,弄好我马上回你。";
if (elapsed < 480_000) return "这件事比预想中多花了一点时间,我还在继续处理。你不用一直盯着,弄好我会直接回你。";
return "我还在处理这件事,确实花了些时间。想看看现在的状态可以发 /status,不想继续等也可以发 /cancel。";
}
const HELP_TEXT = [
"直接用自然语言告诉我你要做什么即可,我会自己判断并处理。",
"兜底命令:",
"/list 列出当前提案",
"/confirm 确认最近待确认的提案",
"/stop 停止正在执行的任务",
"/cancel 取消最近未开始的提案",
"/status 查看运行状态"
].join("\n");