Files
gori-agent/src/core/gateway.ts
T
zenord 91327dc632 feat: share the proposal board across chats
The unfinished-proposal board is now visible to every conversation,
with entries marked scope=own/scope=other (chat type only, no IDs), so
anyone can see what the single worker is busy with. All actions remain
owner-only, and pending reminders still target only the owner.
2026-08-19 17:08:48 +08:00

422 lines
18 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import type { ConversationRuntime } from "../acp/types.js";
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 { Proposal } from "./proposal-store.js";
import type { IncomingMessage, MessageTarget, OutgoingImage, OutgoingImageRef } from "./types.js";
import { readOutboundImage } from "./workspace-images.js";
export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string }
interface ChatTarget {
target: MessageTarget;
adapter: PlatformAdapter;
lastInboundMessageId?: string;
lastInboundAt?: number;
inboundIsGroup: boolean;
lastEventDelivery?: "success" | "failed" | "skipped";
lastEventError?: string;
lastEventAt?: number;
}
interface PendingEvent {
text: string;
images?: OutgoingImageRef[];
}
export interface GatewayOptions {
now?: () => number;
}
// QQ passive-reply windows (with safety margin): 5 minutes for group chats, 60 minutes for C2C.
const GROUP_PASSIVE_WINDOW_MS = 4.5 * 60_000;
const C2C_PASSIVE_WINDOW_MS = 55 * 60_000;
// Reads image files at send time (the worker may have updated them); unreadable files degrade
// to a text note, and adapters without image support get a plain [图片] text line instead.
async function resolveOutboundImages(text: string, refs: OutgoingImageRef[] | undefined, imageCapable: boolean): Promise<{ text: string; images: OutgoingImage[] }> {
if (!refs || refs.length === 0) return { text, images: [] };
const images: OutgoingImage[] = [];
const failures: string[] = [];
for (const ref of refs) {
const image = await readOutboundImage(ref);
if (image) images.push(image);
else failures.push(ref.filename || "image");
}
let finalText = text;
if (failures.length > 0) {
finalText = [finalText, ...failures.map((name) => `(图片 ${name} 发送失败:文件不存在或已失效)`)].filter((part) => part.length > 0).join("\n");
}
if (!imageCapable && images.length > 0) {
finalText = [finalText, ...images.map((image) => `[图片] ${image.filename || "image"}`)].filter((part) => part.length > 0).join("\n");
return { text: finalText, images: [] };
}
return { text: finalText, images };
}
class ReplyStream {
private tail = Promise.resolve();
constructor(
private readonly message: IncomingMessage,
private readonly adapter: PlatformAdapter,
private readonly nextSequence: () => number | undefined
) {}
enqueue(entry: string | PendingEvent, onError?: (error: unknown) => void): Promise<void> {
const request = typeof entry === "string" ? { text: entry } : entry;
this.tail = this.tail.then(() => this.deliver(request)).catch((error: unknown) => {
console.error(`Gateway delivery send failed (platform=${this.message.platform})`);
onError?.(error);
});
return this.tail;
}
private async deliver(request: PendingEvent): Promise<void> {
const { text, images } = await resolveOutboundImages(request.text, request.images, this.adapter.supportsImages === true);
if (text || images.length === 0) await this.send({ text });
for (const image of images) {
try {
await this.send({ text: "", images: [image] });
} catch (error) {
console.error(`Gateway image delivery failed (platform=${this.message.platform})`);
await this.send({ text: `(图片 ${image.filename || "image"} 发送失败)` }).catch(() => undefined);
}
}
}
private send(piece: { text: string; images?: OutgoingImage[] }): Promise<void> {
return this.adapter.sendMessage({
target: {
platform: this.message.platform,
chatId: this.message.chatId,
userId: this.message.userId,
raw: this.message.raw
},
text: piece.text,
replyTo: this.message.messageId,
replySequence: this.nextSequence(),
images: piece.images
});
}
}
export class Gateway {
private readonly chatTargets = new Map<string, ChatTarget>();
private readonly pendingEvents = new Map<string, PendingEvent[]>();
// QQ deduplicates passive replies by msg_id + msg_seq: all sends referencing the same inbound
// message (normal replies, events, redelivered pending events) share one counter per chat+messageId.
private readonly messageSequences = new Map<string, Map<string, number>>();
private readonly now: () => number;
private closed = false;
readonly commandRouter = new CommandRouter();
constructor(
private readonly policy: GatewayPolicy,
private readonly runtime: ConversationRuntime,
options: GatewayOptions = {}
) {
this.now = options.now || Date.now;
}
async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise<GatewayResult> {
const policyError = this.checkPolicy(message);
if (policyError) {
console.log(`Message ignored by policy: ${policyError} (platform=${message.platform})`);
return { ok: true, ignored: true, error: policyError };
}
const chatKey = chatKeyFor(message.platform, message.chatId);
const previous = this.chatTargets.get(chatKey);
this.chatTargets.set(chatKey, {
target: {
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
raw: message.raw
},
adapter,
lastInboundMessageId: message.messageId,
lastInboundAt: this.now(),
inboundIsGroup: Boolean(message.isGroup),
lastEventDelivery: previous?.lastEventDelivery,
lastEventError: previous?.lastEventError,
lastEventAt: previous?.lastEventAt
});
await this.flushPendingEvents(chatKey, message, adapter, options);
const command = this.commandRouter.parse(message.text);
if (command) {
const result = await this.executeCommandResult(command, message);
if (!options.synchronous) await this.replyStream(chatKey, message, adapter).enqueue(result.reply || "");
return result;
}
if (options.synchronous) return this.promptResult(message);
try {
const response = await this.runtime.prompt({
platform: message.platform,
chatId: message.chatId,
userId: message.userId,
text: message.text,
messageId: message.messageId,
attachments: message.attachments
});
await this.replyStream(chatKey, message, adapter).enqueue({ text: response.text, images: response.images });
return { ok: true, reply: response.text };
} catch (error) {
const errorText = error instanceof Error ? error.message : String(error);
const reply = `Agent error: ${errorText}`;
await this.replyStream(chatKey, message, adapter).enqueue(reply);
return { ok: false, error: errorText, reply };
}
}
async sendEvent(chatKey: string, text: string, images?: OutgoingImageRef[]): Promise<void> {
if (this.closed) return;
const entry = this.chatTargets.get(chatKey);
if (!entry) return;
const windowMs = entry.inboundIsGroup ? GROUP_PASSIVE_WINDOW_MS : C2C_PASSIVE_WINDOW_MS;
const freshInbound = entry.lastInboundMessageId
&& entry.lastInboundAt !== undefined
&& this.now() - entry.lastInboundAt <= windowMs;
if (!freshInbound) {
// No fresh passive window and proactive messages may be unauthorized: skip and remind on the next inbound.
this.recordEventDelivery(chatKey, "skipped", "no fresh passive window");
this.queuePendingEvent(chatKey, { text, images });
console.log(`Gateway event skipped: no fresh passive window (platform=${entry.target.platform})`);
return;
}
try {
const resolved = await resolveOutboundImages(text, images, entry.adapter.supportsImages === true);
if (resolved.text || resolved.images.length === 0) {
await entry.adapter.sendMessage({
target: entry.target,
text: resolved.text,
replyTo: entry.lastInboundMessageId,
replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId)
});
}
for (const image of resolved.images) {
try {
await entry.adapter.sendMessage({
target: entry.target,
text: "",
replyTo: entry.lastInboundMessageId,
replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId),
images: [image]
});
} catch (imageError) {
// Image upload/send failures never block the event text; degrade to a text note.
console.error(`Gateway event image send failed (platform=${entry.target.platform})`);
await entry.adapter.sendMessage({
target: entry.target,
text: `(图片 ${image.filename || "image"} 发送失败)`,
replyTo: entry.lastInboundMessageId,
replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId)
}).catch(() => undefined);
}
}
this.recordEventDelivery(chatKey, "success");
} catch (error) {
this.recordEventDelivery(chatKey, "failed", safeEventError(error));
this.queuePendingEvent(chatKey, { text, images });
console.error(`Gateway event send failed (platform=${entry.target.platform})`);
}
}
private recordEventDelivery(chatKey: string, delivery: "success" | "failed" | "skipped", error?: string): void {
const entry = this.chatTargets.get(chatKey);
if (!entry) return;
entry.lastEventDelivery = delivery;
entry.lastEventError = error;
entry.lastEventAt = this.now();
}
private queuePendingEvent(chatKey: string, event: PendingEvent): void {
const pending = this.pendingEvents.get(chatKey) || [];
pending.push(event);
this.pendingEvents.set(chatKey, pending);
}
private async flushPendingEvents(chatKey: string, message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean }): Promise<void> {
if (options.synchronous) return;
const pending = this.pendingEvents.get(chatKey);
if (!pending || pending.length === 0) return;
this.pendingEvents.delete(chatKey);
const stream = this.replyStream(chatKey, message, adapter);
for (const event of pending) {
// Images are re-read from disk at redelivery time (the worker may have updated them).
await stream.enqueue(event, (error) => {
// Redelivery failed: keep the event queued for the next inbound and record the failure.
this.recordEventDelivery(chatKey, "failed", safeEventError(error));
this.queuePendingEvent(chatKey, event);
});
}
}
private replyStream(chatKey: string, message: IncomingMessage, adapter: PlatformAdapter): ReplyStream {
return new ReplyStream(message, adapter, () => this.nextReplySequence(chatKey, message.messageId));
}
private nextReplySequence(chatKey: string, messageId: string | undefined): number | undefined {
if (!messageId) return undefined;
let sequences = this.messageSequences.get(chatKey);
if (!sequences) {
sequences = new Map();
this.messageSequences.set(chatKey, sequences);
}
const next = (sequences.get(messageId) || 0) + 1;
// Re-set to refresh insertion order, then bound memory to the most recent message ids per chat.
sequences.delete(messageId);
sequences.set(messageId, next);
while (sequences.size > 8) sequences.delete(sequences.keys().next().value!);
return next;
}
private eventStatus(chatKey: string): Record<string, string> {
const entry = this.chatTargets.get(chatKey);
return {
lastEventDelivery: entry?.lastEventDelivery || "none",
lastEventError: entry?.lastEventError || "none",
lastEventAt: entry?.lastEventAt ? new Date(entry.lastEventAt).toISOString() : "never"
};
}
stats(): ReturnType<ConversationRuntime["stats"]> {
return this.runtime.stats();
}
shutdown(): void {
this.closed = true;
}
private async executeCommandResult(command: ParsedCommand, message: IncomingMessage): Promise<GatewayResult> {
try {
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,
text: message.text,
messageId: message.messageId,
attachments: message.attachments
})).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) {
case "help": return HELP_TEXT;
case "status": {
const status = this.runtime.status(message.platform, message.chatId, message.userId);
const combined = { ...status, ...this.eventStatus(chatKeyFor(message.platform, message.chatId)) };
return `OK\n${Object.entries(combined).map(([key, value]) => `${key}=${value}`).join("\n")}`;
}
case "list": {
if (!this.runtime.listProposals) return "当前运行时不支持列出提案。";
const chatKey = chatKeyFor(message.platform, message.chatId);
const owned = (proposal: Proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === message.userId;
return renderProposalPanel(this.runtime.listProposals(message.platform, message.chatId, message.userId), owned);
}
case "confirm": {
if (!this.runtime.confirm) return "当前运行时不支持确认操作。";
return await this.runtime.confirm(message.platform, message.chatId, message.userId) ? "已确认,加入队列。" : "当前没有待确认的提案。";
}
case "finish": {
if (!this.runtime.finish) return "当前运行时不支持结束操作。";
return await this.runtime.finish(message.platform, message.chatId, message.userId) ? "已结束该任务。" : "当前没有待结束的任务。";
}
case "stop": {
if (!this.runtime.stop) return "当前运行时不支持停止操作。";
return await this.runtime.stop(message.platform, message.chatId, message.userId) ? "已停止当前任务,任务转为待确认。" : "当前没有正在执行的任务。";
}
case "cancel":
return await this.runtime.cancel(message.platform, message.chatId, message.userId) ? "已取消最近的提案。" : "没有可取消的提案。";
}
}
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";
if (this.policy.requireMentionInGroup && message.isGroup && !message.mentionsBot) return "Mention required in group chat";
return undefined;
}
}
function renderProposalPanel(proposals: Proposal[], owned: (proposal: Proposal) => boolean): string {
if (proposals.length === 0) return "当前没有提案。";
// The board is globally visible; ownership is annotated without exposing chat/user IDs.
const scope = (proposal: Proposal) => owned(proposal) ? "(你的)" : "(其他成员)";
const pending = proposals.filter((proposal) => proposal.status === "pending");
const working = proposals.filter((proposal) => proposal.status === "working");
const queued = proposals.filter((proposal) => proposal.status === "queued");
const proposed = proposals.filter((proposal) => proposal.status === "proposed");
const finished = proposals.filter((proposal) => proposal.status === "finished")
.sort((left, right) => (right.finishedAt ?? right.updatedAt) - (left.finishedAt ?? left.updatedAt))
.slice(0, 3);
const sections: string[] = [];
if (pending.length > 0) {
sections.push("待确认(等你处理):");
for (const proposal of pending) {
const card = proposal.pending;
sections.push(`- ${proposal.id}「${proposal.title}」${scope(proposal)}${card ? `:${card.summary}${card.question ? `(问你:${card.question})` : ""}` : ""}(说 /finish 结束,或直接说要求继续改)`);
}
}
if (working.length > 0) {
sections.push("进行中:");
for (const proposal of working) sections.push(`- ${proposal.id}「${proposal.title}」${scope(proposal)}正在执行(/stop 可中止)`);
}
if (queued.length > 0) {
sections.push("排队中:");
for (const proposal of queued) sections.push(`- ${proposal.id}「${proposal.title}」${scope(proposal)}`);
}
if (proposed.length > 0) {
sections.push("未确认(proposed):");
for (const proposal of proposed) sections.push(`- ${proposal.id}「${proposal.title}」${scope(proposal)}(/confirm 确认后排队)`);
}
if (finished.length > 0) {
sections.push("最近结束:");
for (const proposal of finished) {
const state = proposal.finishKind === "cancelled" ? "已取消" : "已完成";
sections.push(`- ${proposal.id}「${proposal.title}」${scope(proposal)}${state}${proposal.finishNote ? `(${proposal.finishNote})` : ""}`);
}
}
return sections.join("\n");
}
function safeEventError(error: unknown): string {
const message = error instanceof Error ? error.message : "";
const http = /HTTP (\d{3})/.exec(message);
const code = /code=(\d+)/.exec(message);
const parts = [http ? `HTTP ${http[1]}` : "", code ? `QQ code ${code[1]}` : ""].filter(Boolean);
if (parts.length > 0) return parts.join(" ");
return error instanceof Error ? error.name : "error";
}
const HELP_TEXT = [
"直接用自然语言告诉我你要做什么即可,我会自己判断并处理。",
"兜底命令:",
"/list 列出当前提案",
"/confirm 确认最近待确认的提案",
"/finish 结束最近待确认的任务",
"/stop 停止正在执行的任务",
"/cancel 取消最近未开始或待确认的提案",
"/status 查看运行状态"
].join("\n");