feat: delay long-running task notices

This commit is contained in:
zenord
2026-08-17 15:35:18 +08:00
parent 2795567f8c
commit 680a449cd8
6 changed files with 589 additions and 164 deletions
+203 -35
View File
@@ -7,18 +7,75 @@ import type { IncomingMessage } 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>;
}
class ReplyStream {
private sequence = 0;
private tail = Promise.resolve();
constructor(private readonly message: IncomingMessage, private readonly adapter: PlatformAdapter) {}
enqueue(text: string): Promise<void> {
const replySequence = ++this.sequence;
this.tail = this.tail.then(() => this.adapter.sendMessage({
target: {
platform: this.message.platform,
chatId: this.message.chatId,
userId: this.message.userId,
raw: this.message.raw
},
text,
replyTo: this.message.messageId,
replySequence
})).catch(() => {
console.error(`Gateway delivery send failed (platform=${this.message.platform})`);
});
return this.tail;
}
}
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 closed = false;
readonly commandRouter = new CommandRouter();
constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime) {}
constructor(
private readonly policy: GatewayPolicy,
private readonly runtime: ConversationRuntime,
options: GatewayOptions = {}
) {
this.now = options.now || Date.now;
this.schedule = options.schedule || defaultSchedule;
}
async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise<GatewayResult> {
const policyError = this.checkPolicy(message);
@@ -29,33 +86,92 @@ export class Gateway {
const command = this.commandRouter.parse(message.text);
if (command?.kind === "help" || command?.kind === "status" || command?.kind === "cancel") {
return this.reply(message, adapter, await this.executeCommand(command, message), options);
return this.executeAndReply(command, message, adapter, options);
}
const chatKey = chatKeyFor(message.platform, message.chatId);
const existing = this.chatWork.get(chatKey);
const wasBusy = Boolean(existing?.running || existing?.queued);
const acknowledgement = !command && !options.synchronous
? this.sendAcknowledgement(message, adapter, wasBusy)
: Promise.resolve();
return this.withChatLock(chatKey, async () => {
await acknowledgement;
const replySequence = !command && !options.synchronous ? 2 : 1;
if (command || options.synchronous) {
const result = await this.withChatLock(chatKey, () => this.executeSynchronous(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, message);
const locked = await this.withChatLock(chatKey, async () => {
inbound.phase = "running";
if (inbound.queuedNoticeSent) void inbound.stream.enqueue("轮到这条了,我开始处理。");
try {
const reply = command ? await this.executeCommand(command, message) : (await this.runtime.prompt({
platform: message.platform, chatId: message.chatId, userId: message.userId, text: message.text, messageId: message.messageId
const reply = (await this.runtime.prompt({
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
text: message.text,
messageId: message.messageId
})).text;
return this.reply(message, adapter, reply, options, replySequence);
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}`;
if (!options.synchronous) await this.send(message, adapter, reply, replySequence);
return { ok: false, error: errorText, reply };
this.finishInbound(inbound);
return { result: { ok: false, error: errorText, reply }, delivery: inbound.stream.enqueue(reply) };
}
});
await locked.delivery;
return locked.result;
}
stats(): ReturnType<ConversationRuntime["stats"]> & { lockedChats: number } { return { ...this.runtime.stats(), lockedChats: this.locks.size }; }
stats(): ReturnType<ConversationRuntime["stats"]> & { lockedChats: number } {
return { ...this.runtime.stats(), lockedChats: this.locks.size };
}
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> {
try {
const reply = command ? await this.executeCommand(command, message) : (await this.runtime.prompt({
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
text: message.text,
messageId: message.messageId
})).text;
return { ok: true, reply };
} catch (error) {
const errorText = error instanceof Error ? error.message : String(error);
return { ok: false, error: errorText, reply: `Agent error: ${errorText}` };
}
}
private async executeCommand(command: ParsedCommand, message: IncomingMessage): Promise<string> {
switch (command.kind) {
@@ -67,7 +183,8 @@ export class Gateway {
...this.runtime.status(message.platform, message.chatId),
gatewayRunning: Boolean(work?.running),
queued: work?.queued || 0,
gatewayRunningSeconds: work?.startedAt ? Math.max(0, Math.floor((Date.now() - work.startedAt) / 1000)) : 0
gatewayRunningSeconds: work?.startedAt !== undefined
? Math.max(0, Math.floor((this.now() - work.startedAt) / 1000)) : 0
};
return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`;
}
@@ -83,26 +200,55 @@ export class Gateway {
return await this.runtime.cancel(message.platform, message.chatId) ? "Cancellation requested." : "No active turn to cancel.";
}
private async reply(message: IncomingMessage, adapter: PlatformAdapter, reply: string, options: { synchronous?: boolean }, replySequence = 1): Promise<GatewayResult> {
if (!options.synchronous) await this.send(message, adapter, reply, replySequence);
return { ok: true, reply };
private armTimer(inbound: AsyncInbound, deadline: number, message: IncomingMessage): 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, message);
return;
}
if (!inbound.initialNoticeSent) {
inbound.initialNoticeSent = true;
inbound.queuedNoticeSent = inbound.phase === "queued";
void inbound.stream.enqueue(inbound.phase === "queued"
? "前面还有任务,这条还在排队,轮到后我马上处理。"
: "这条还在处理,完成后我会直接回复。");
} else {
void inbound.stream.enqueue(this.reminderText(inbound, message));
}
this.armTimer(inbound, nextReminderDeadline(inbound.receivedAt, now), message);
}, Math.max(0, deadline - this.now()));
}
private send(message: IncomingMessage, adapter: PlatformAdapter, text: string, replySequence: number): Promise<void> {
return adapter.sendMessage({
target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw },
text,
replyTo: message.messageId,
replySequence
});
private reminderText(inbound: AsyncInbound, message: IncomingMessage): string {
let idleSeconds: number | undefined;
try {
const value = this.runtime.status(message.platform, message.chatId).idleSeconds;
if (typeof value === "number" && Number.isFinite(value) && value >= 0) idleSeconds = Math.floor(value);
} catch {
// A status sampling failure must not affect the active turn.
}
const activity = activityText(idleSeconds);
if (inbound.phase === "queued") {
return activity
? `前面的任务还在处理,${activity};这条仍在排队,轮到后我马上处理。`
: "前面的任务还在处理,这条仍在排队,轮到后我马上处理。";
}
return activity
? `这条还在处理中,${activity};完成后我会直接回复。`
: "这条还在处理中,完成后我会直接回复。";
}
private async sendAcknowledgement(message: IncomingMessage, adapter: PlatformAdapter, queued: boolean): Promise<void> {
const text = queued
? "已收到,已排队。可随时发送 /status 查看状态。"
: "已收到,正在处理。可随时发送 /status 查看状态。";
try { await this.send(message, adapter, text, 1); }
catch { console.error(`Gateway acknowledgement send failed (platform=${message.platform})`); }
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 {
@@ -124,8 +270,10 @@ export class Gateway {
await previous;
state.queued--;
state.running = true;
state.startedAt = Date.now();
try { return await fn(); } finally {
state.startedAt = this.now();
try {
return await fn();
} finally {
state.running = false;
state.startedAt = undefined;
release();
@@ -136,3 +284,23 @@ export class Gateway {
}
}
}
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;
return receivedAt + 180_000 + (Math.floor((elapsed - 180_000) / 300_000) + 1) * 300_000;
}
function activityText(idleSeconds: number | undefined): string | undefined {
if (idleSeconds === undefined) return undefined;
if (idleSeconds < 30) return "刚刚还有进展";
if (idleSeconds < 120) return `最近 ${idleSeconds} 秒没新进展`;
return `最近 ${Math.floor(idleSeconds / 60)} 分钟没新进展`;
}