feat: send workspace images to QQ chats
Worker results may report image attachments stored inside the workspace; the runtime validates them (containment, png/jpg magic, size) and delivers them with the pending event. Assistants gain a send_image action so users can ask for an image later. QQ uploads via /files with base64 file_data and sends msg_type 7 rich media, sharing the same msg_id/msg_seq counter as text replies; non-image adapters flatten images to text.
This commit is contained in:
@@ -4,7 +4,8 @@ import path from "node:path";
|
||||
import type { AcpConfig, PermissionPolicy } from "../config.js";
|
||||
import { chatKeyFor, type AssistantBinding, type DurableSessionStore } from "../core/durable-session-store.js";
|
||||
import type { Proposal, ProposalPending, ProposalStore } from "../core/proposal-store.js";
|
||||
import type { IncomingAttachment } from "../core/types.js";
|
||||
import type { IncomingAttachment, OutgoingImageRef } from "../core/types.js";
|
||||
import { MAX_OUTBOUND_IMAGES, validateWorkspaceImage } from "../core/workspace-images.js";
|
||||
import type { ResolvedBot } from "../roles/role-registry.js";
|
||||
import type { ConversationRequest, ConversationResponse, ConversationRuntime, RuntimeEventSink, RuntimeStats } from "./types.js";
|
||||
import type { AcpPromptContent } from "./client.js";
|
||||
@@ -20,6 +21,7 @@ export type AssistantAction =
|
||||
| { type: "adjust_proposal"; id?: string; title?: string; goal?: string; steps?: string[] }
|
||||
| { type: "follow_up"; id?: string; instruction: string }
|
||||
| { type: "finish"; id?: string; note?: string }
|
||||
| { type: "send_image"; path: string }
|
||||
| { type: "start_next" }
|
||||
| { type: "cancel"; id?: string }
|
||||
| { type: "stop" };
|
||||
@@ -36,6 +38,7 @@ export interface ParsedWorkerResult {
|
||||
summary: string;
|
||||
question?: string;
|
||||
workspaceDirty?: boolean;
|
||||
attachments?: { path: string; mimeType?: string }[];
|
||||
}
|
||||
|
||||
interface ActiveWorker {
|
||||
@@ -111,8 +114,8 @@ export class AssistantManager implements ConversationRuntime {
|
||||
async prompt(request: ConversationRequest): Promise<ConversationResponse> {
|
||||
if (this.shuttingDown) throw new Error("ACP runtime is shutting down");
|
||||
const conversationKey = assistantKeyFor(chatKeyFor(request.platform, request.chatId), request.userId);
|
||||
const reply = await this.enqueueChat(conversationKey, () => this.runAssistantTurn(conversationKey, request));
|
||||
return { text: reply, botId: this.bot.id, agentId: this.bot.agent.id, mode: "assistant" };
|
||||
const result = await this.enqueueChat(conversationKey, () => this.runAssistantTurn(conversationKey, request));
|
||||
return { text: result.text, images: result.images, botId: this.bot.id, agentId: this.bot.agent.id, mode: "assistant" };
|
||||
}
|
||||
|
||||
async cancel(platform: string, chatId: string, userId: string): Promise<boolean> {
|
||||
@@ -226,7 +229,7 @@ export class AssistantManager implements ConversationRuntime {
|
||||
await Promise.allSettled([...this.chatChains.values()]);
|
||||
}
|
||||
|
||||
private async runAssistantTurn(conversationKey: string, request: ConversationRequest): Promise<string> {
|
||||
private async runAssistantTurn(conversationKey: string, request: ConversationRequest): Promise<{ text: string; images?: OutgoingImageRef[] }> {
|
||||
const worker = await this.acquireAssistantWorker(conversationKey);
|
||||
try {
|
||||
const reply = await worker.prompt(this.assistantPromptContent(worker, request));
|
||||
@@ -234,10 +237,10 @@ export class AssistantManager implements ConversationRuntime {
|
||||
if (worker.isolationViolated) throw new Error("Assistant session attempted forbidden tool activity");
|
||||
await this.store.touchBinding(conversationKey);
|
||||
const chatKey = chatKeyFor(request.platform, request.chatId);
|
||||
const corrections = await this.executeActions(chatKey, request, parsed.actions);
|
||||
const { corrections, images } = await this.executeActions(chatKey, request, parsed.actions);
|
||||
const reminders = this.pendingReminders(chatKey, request.userId, parsed.reply);
|
||||
const extras = [...corrections, ...reminders];
|
||||
return extras.length > 0 ? `${parsed.reply}\n${extras.join("\n")}` : parsed.reply;
|
||||
return { text: extras.length > 0 ? `${parsed.reply}\n${extras.join("\n")}` : parsed.reply, images };
|
||||
} catch (error) {
|
||||
if (worker.isolationViolated) await this.store.deleteBinding(conversationKey).catch(() => undefined);
|
||||
throw error;
|
||||
@@ -269,8 +272,9 @@ export class AssistantManager implements ConversationRuntime {
|
||||
return parsed;
|
||||
}
|
||||
|
||||
private async executeActions(chatKey: string, request: ConversationRequest, actions: AssistantAction[]): Promise<string[]> {
|
||||
private async executeActions(chatKey: string, request: ConversationRequest, actions: AssistantAction[]): Promise<{ corrections: string[]; images?: OutgoingImageRef[] }> {
|
||||
const corrections: string[] = [];
|
||||
const images: OutgoingImageRef[] = [];
|
||||
for (const action of actions) {
|
||||
if (action.type === "create_proposal") {
|
||||
await this.proposals.create({
|
||||
@@ -292,6 +296,10 @@ export class AssistantManager implements ConversationRuntime {
|
||||
} else if (action.type === "finish") {
|
||||
const ok = await this.scheduler(() => this.finishLocked(chatKey, request.userId, action.id, action.note));
|
||||
if (!ok) corrections.push("没有可结束的任务;只有待确认的任务可以 finish,且你只能操作自己发起的任务。");
|
||||
} else if (action.type === "send_image") {
|
||||
const image = await validateWorkspaceImage(this.bot.workspace, action.path);
|
||||
if (image) images.push({ path: image.path, mimeType: image.mimeType, filename: image.filename });
|
||||
else corrections.push("这张图发不出去:路径不在工作区内,或者不是有效的 png/jpg 图片。");
|
||||
} else if (action.type === "start_next") {
|
||||
const started = await this.scheduler(() => this.tryStartLocked(chatKey, request.userId));
|
||||
if (!started) corrections.push(this.startNextCorrection(chatKey, request.userId));
|
||||
@@ -303,7 +311,7 @@ export class AssistantManager implements ConversationRuntime {
|
||||
if (!ok) corrections.push("没有可停止的任务;你只能操作自己发起的任务。");
|
||||
}
|
||||
}
|
||||
return corrections;
|
||||
return { corrections, images: images.length > 0 ? images.slice(0, MAX_OUTBOUND_IMAGES) : undefined };
|
||||
}
|
||||
|
||||
private startNextCorrection(chatKey: string, userId: string): string {
|
||||
@@ -517,11 +525,23 @@ export class AssistantManager implements ConversationRuntime {
|
||||
if (!proposal || proposal.status !== "working") { active.settled = true; return false; }
|
||||
active.settled = true;
|
||||
this.active = undefined;
|
||||
const attachments: { path: string; mimeType?: string }[] = [];
|
||||
const droppedAttachments: string[] = [];
|
||||
for (const ref of result.attachments || []) {
|
||||
const image = await validateWorkspaceImage(this.bot.workspace, ref.path);
|
||||
if (image) attachments.push({ path: image.path, mimeType: image.mimeType });
|
||||
else droppedAttachments.push(ref.path);
|
||||
}
|
||||
if (droppedAttachments.length > 0) {
|
||||
console.error(`Worker reported ${droppedAttachments.length} invalid attachment(s) for proposal ${active.proposalId}; dropped (outside workspace or not a valid png/jpg)`);
|
||||
}
|
||||
const pending: ProposalPending = {
|
||||
summary: result.summary,
|
||||
question: result.question,
|
||||
workspaceDirty: result.workspaceDirty,
|
||||
receivedAt: Date.now()
|
||||
receivedAt: Date.now(),
|
||||
attachments: attachments.length > 0 ? attachments : undefined,
|
||||
droppedAttachments: droppedAttachments.length > 0 ? droppedAttachments : undefined
|
||||
};
|
||||
await this.proposals.update(active.proposalId, {
|
||||
status: "pending",
|
||||
@@ -578,11 +598,11 @@ export class AssistantManager implements ConversationRuntime {
|
||||
if (this.shuttingDown) return;
|
||||
try {
|
||||
const reply = await this.runAssistantEvent(conversationKey, proposal.id);
|
||||
if (reply) await sink(chatKey, reply);
|
||||
if (reply) await sink(chatKey, reply, pendingImageRefs(this.proposals.get(proposalId)));
|
||||
} catch (error) {
|
||||
console.error(`Assistant event notification failed (${error instanceof Error ? error.name : "unknown error"})`);
|
||||
const latest = this.proposals.get(proposalId);
|
||||
if (latest?.pending) await sink(chatKey, fallbackEventText(latest)).catch(() => undefined);
|
||||
if (latest?.pending) await sink(chatKey, fallbackEventText(latest), pendingImageRefs(latest)).catch(() => undefined);
|
||||
}
|
||||
}).catch(() => undefined);
|
||||
}
|
||||
@@ -844,6 +864,10 @@ function parseAssistantAction(value: unknown): AssistantAction | undefined {
|
||||
if ((value.id !== undefined && typeof value.id !== "string") || (value.note !== undefined && typeof value.note !== "string")) return undefined;
|
||||
return { type: "finish", id: value.id as string | undefined, note: value.note as string | undefined };
|
||||
}
|
||||
if (value.type === "send_image") {
|
||||
if (typeof value.path !== "string" || value.path.length === 0) return undefined;
|
||||
return { type: "send_image", path: value.path };
|
||||
}
|
||||
if (value.type === "start_next") return { type: "start_next" };
|
||||
if (value.type === "cancel") {
|
||||
if (value.id !== undefined && typeof value.id !== "string") return undefined;
|
||||
@@ -862,16 +886,27 @@ export function parseWorkerResult(text: string): ParsedWorkerResult | undefined
|
||||
if (value.status !== "PENDING") return undefined;
|
||||
if ((value.question !== undefined && typeof value.question !== "string")
|
||||
|| (value.workspaceDirty !== undefined && typeof value.workspaceDirty !== "boolean")) return undefined;
|
||||
let attachments: { path: string; mimeType?: string }[] | undefined;
|
||||
if (value.attachments !== undefined) {
|
||||
if (!Array.isArray(value.attachments) || value.attachments.length > MAX_OUTBOUND_IMAGES) return undefined;
|
||||
attachments = [];
|
||||
for (const item of value.attachments) {
|
||||
if (!isRecord(item) || typeof item.path !== "string" || item.path.length === 0
|
||||
|| (item.mimeType !== undefined && typeof item.mimeType !== "string")) return undefined;
|
||||
attachments.push({ path: item.path, mimeType: item.mimeType as string | undefined });
|
||||
}
|
||||
}
|
||||
return {
|
||||
status: "PENDING",
|
||||
summary: value.summary,
|
||||
question: value.question as string | undefined,
|
||||
workspaceDirty: value.workspaceDirty as boolean | undefined
|
||||
workspaceDirty: value.workspaceDirty as boolean | undefined,
|
||||
attachments
|
||||
};
|
||||
}
|
||||
|
||||
function withWorkerResultProtocol(input: string | AcpPromptContent[]): string | AcpPromptContent[] {
|
||||
const instruction = "When this turn is finished, end your response with exactly one hidden worker result envelope: <GORI_WORKER_RESULT_V2>{\"status\":\"PENDING\",\"summary\":\"...\"}</GORI_WORKER_RESULT_V2>. The summary is a short user-readable description of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing, and set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. PENDING is the only status; do not emit any other status value or any text after the envelope.";
|
||||
const instruction = "When this turn is finished, end your response with exactly one hidden worker result envelope: <GORI_WORKER_RESULT_V2>{\"status\":\"PENDING\",\"summary\":\"...\"}</GORI_WORKER_RESULT_V2>. The summary is a short user-readable description of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing, and set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. When you produced image files the user should see (png/jpg only), save them inside the workspace (prefer .gori-outbox/) and report up to 3 of them as \"attachments\": [{\"path\":\"relative/or/absolute/path\"}]; never report paths outside the workspace. PENDING is the only status; do not emit any other status value or any text after the envelope.";
|
||||
if (typeof input === "string") return `${input}\n\n${instruction}`;
|
||||
return [...input, { type: "text" as const, text: `\n\n${instruction}` }];
|
||||
}
|
||||
@@ -908,18 +943,38 @@ function workerFollowUpPrompt(instruction: string, question?: string): string {
|
||||
|
||||
function assistantEventPrompt(proposal: Proposal): string {
|
||||
const pending = proposal.pending!;
|
||||
const attachments = pending.attachments?.length
|
||||
? `Image attachments produced by the worker (already sent to the user after your message): ${pending.attachments.map((attachment) => attachment.path).join(", ")}`
|
||||
: "Image attachments: none";
|
||||
const dropped = pending.droppedAttachments?.length
|
||||
? `Note: ${pending.droppedAttachments.length} reported attachment(s) were invalid (outside the workspace or not png/jpg) and were dropped; mention this to the user.`
|
||||
: "";
|
||||
return [
|
||||
"[Internal event from gori-agent. This is not a user message.]",
|
||||
`The worker for proposal "${proposal.title}" (id ${proposal.id}) reported its result and the proposal is now waiting for the user's decision.`,
|
||||
`Summary: ${pending.summary}`,
|
||||
pending.question ? `Question for the user: ${pending.question}` : "Question for the user: none",
|
||||
attachments,
|
||||
...(dropped ? [dropped] : []),
|
||||
"Tell the user, in your own words, what happened and what they can do next (finish to close the task, follow up with new instructions to continue it, or stop). End with the usual GORI_ASSISTANT_ACTION_V2 envelope; its actions array MUST be empty because actions are ignored for internal events."
|
||||
].join("\n");
|
||||
}
|
||||
|
||||
function fallbackEventText(proposal: Proposal): string {
|
||||
const pending = proposal.pending!;
|
||||
return `任务「${proposal.title}」等你确认。摘要:${pending.summary}${pending.question ? ` 问题:${pending.question}` : ""} 说 finish 结束,或直接说要求继续改。`;
|
||||
const attachments = pending.attachments?.length ? `(附 ${pending.attachments.length} 张图片)` : "";
|
||||
const dropped = pending.droppedAttachments?.length ? `(${pending.droppedAttachments.length} 个附件校验失败被丢弃)` : "";
|
||||
return `任务「${proposal.title}」等你确认。摘要:${pending.summary}${pending.question ? ` 问题:${pending.question}` : ""}${attachments}${dropped} 说 finish 结束,或直接说要求继续改。`;
|
||||
}
|
||||
|
||||
function pendingImageRefs(proposal: Proposal | undefined): OutgoingImageRef[] | undefined {
|
||||
const attachments = proposal?.pending?.attachments;
|
||||
if (!attachments || attachments.length === 0) return undefined;
|
||||
return attachments.map((attachment) => ({
|
||||
path: attachment.path,
|
||||
mimeType: attachment.mimeType,
|
||||
filename: attachment.path.split("/").pop()
|
||||
}));
|
||||
}
|
||||
|
||||
function isRecord(value: unknown): value is Record<string, unknown> {
|
||||
|
||||
+3
-2
@@ -1,6 +1,6 @@
|
||||
import type { PermissionPolicy } from "../config.js";
|
||||
import type { Proposal } from "../core/proposal-store.js";
|
||||
import type { IncomingAttachment } from "../core/types.js";
|
||||
import type { IncomingAttachment, OutgoingImageRef } from "../core/types.js";
|
||||
|
||||
export interface AcpBackendSpec {
|
||||
id: string;
|
||||
@@ -23,6 +23,7 @@ export interface ConversationResponse {
|
||||
botId: string;
|
||||
agentId: string;
|
||||
mode?: "assistant";
|
||||
images?: OutgoingImageRef[];
|
||||
}
|
||||
|
||||
export interface RuntimeStats {
|
||||
@@ -32,7 +33,7 @@ export interface RuntimeStats {
|
||||
persistedBindings: number;
|
||||
}
|
||||
|
||||
export type RuntimeEventSink = (chatKey: string, text: string) => Promise<void>;
|
||||
export type RuntimeEventSink = (chatKey: string, text: string, images?: OutgoingImageRef[]) => Promise<void>;
|
||||
|
||||
export interface ConversationRuntime {
|
||||
prompt(request: ConversationRequest): Promise<ConversationResponse>;
|
||||
|
||||
@@ -2,6 +2,9 @@ import type { OutgoingMessage, WebhookRequestContext, WebhookResponse } from "./
|
||||
|
||||
export interface PlatformAdapter {
|
||||
name: string;
|
||||
// True when the adapter can deliver OutgoingMessage.images (e.g. QQ msg_type 7);
|
||||
// gateways degrade images to text notes for adapters without this flag.
|
||||
supportsImages?: boolean;
|
||||
handleWebhook(context: WebhookRequestContext): Promise<WebhookResponse>;
|
||||
sendMessage(message: OutgoingMessage): Promise<void>;
|
||||
}
|
||||
|
||||
+99
-29
@@ -4,7 +4,8 @@ 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 } from "./types.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 }
|
||||
|
||||
@@ -19,6 +20,11 @@ interface ChatTarget {
|
||||
lastEventAt?: number;
|
||||
}
|
||||
|
||||
interface PendingEvent {
|
||||
text: string;
|
||||
images?: OutgoingImageRef[];
|
||||
}
|
||||
|
||||
export interface GatewayOptions {
|
||||
now?: () => number;
|
||||
}
|
||||
@@ -27,6 +33,28 @@ export interface GatewayOptions {
|
||||
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();
|
||||
|
||||
@@ -36,29 +64,47 @@ class ReplyStream {
|
||||
private readonly nextSequence: () => number | undefined
|
||||
) {}
|
||||
|
||||
enqueue(text: string, onError?: (error: unknown) => void): Promise<void> {
|
||||
const replySequence = this.nextSequence();
|
||||
this.tail = this.tail.then(() => this.adapter.sendMessage({
|
||||
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,
|
||||
text: piece.text,
|
||||
replyTo: this.message.messageId,
|
||||
replySequence
|
||||
})).catch((error: unknown) => {
|
||||
console.error(`Gateway delivery send failed (platform=${this.message.platform})`);
|
||||
onError?.(error);
|
||||
replySequence: this.nextSequence(),
|
||||
images: piece.images
|
||||
});
|
||||
return this.tail;
|
||||
}
|
||||
}
|
||||
|
||||
export class Gateway {
|
||||
private readonly chatTargets = new Map<string, ChatTarget>();
|
||||
private readonly pendingEvents = new Map<string, string[]>();
|
||||
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>>();
|
||||
@@ -110,16 +156,16 @@ export class Gateway {
|
||||
if (options.synchronous) return this.promptResult(message);
|
||||
|
||||
try {
|
||||
const reply = (await this.runtime.prompt({
|
||||
const response = await this.runtime.prompt({
|
||||
platform: message.platform,
|
||||
chatId: message.chatId,
|
||||
userId: message.userId,
|
||||
text: message.text,
|
||||
messageId: message.messageId,
|
||||
attachments: message.attachments
|
||||
})).text;
|
||||
await this.replyStream(chatKey, message, adapter).enqueue(reply);
|
||||
return { ok: true, reply };
|
||||
});
|
||||
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}`;
|
||||
@@ -128,7 +174,7 @@ export class Gateway {
|
||||
}
|
||||
}
|
||||
|
||||
async sendEvent(chatKey: string, text: string): Promise<void> {
|
||||
async sendEvent(chatKey: string, text: string, images?: OutgoingImageRef[]): Promise<void> {
|
||||
if (this.closed) return;
|
||||
const entry = this.chatTargets.get(chatKey);
|
||||
if (!entry) return;
|
||||
@@ -139,21 +185,44 @@ export class Gateway {
|
||||
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);
|
||||
this.queuePendingEvent(chatKey, { text, images });
|
||||
console.log(`Gateway event skipped: no fresh passive window (platform=${entry.target.platform})`);
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await entry.adapter.sendMessage({
|
||||
target: entry.target,
|
||||
text,
|
||||
replyTo: entry.lastInboundMessageId,
|
||||
replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId)
|
||||
});
|
||||
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);
|
||||
this.queuePendingEvent(chatKey, { text, images });
|
||||
console.error(`Gateway event send failed (platform=${entry.target.platform})`);
|
||||
}
|
||||
}
|
||||
@@ -166,9 +235,9 @@ export class Gateway {
|
||||
entry.lastEventAt = this.now();
|
||||
}
|
||||
|
||||
private queuePendingEvent(chatKey: string, text: string): void {
|
||||
private queuePendingEvent(chatKey: string, event: PendingEvent): void {
|
||||
const pending = this.pendingEvents.get(chatKey) || [];
|
||||
pending.push(text);
|
||||
pending.push(event);
|
||||
this.pendingEvents.set(chatKey, pending);
|
||||
}
|
||||
|
||||
@@ -178,11 +247,12 @@ export class Gateway {
|
||||
if (!pending || pending.length === 0) return;
|
||||
this.pendingEvents.delete(chatKey);
|
||||
const stream = this.replyStream(chatKey, message, adapter);
|
||||
for (const text of pending) {
|
||||
await stream.enqueue(text, (error) => {
|
||||
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, text);
|
||||
this.queuePendingEvent(chatKey, event);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -19,11 +19,18 @@ export type ProposalStatus =
|
||||
| "pending"
|
||||
| "finished";
|
||||
|
||||
export interface ProposalPendingAttachment {
|
||||
path: string;
|
||||
mimeType?: string;
|
||||
}
|
||||
|
||||
export interface ProposalPending {
|
||||
summary: string;
|
||||
question?: string;
|
||||
workspaceDirty?: boolean;
|
||||
receivedAt: number;
|
||||
attachments?: ProposalPendingAttachment[];
|
||||
droppedAttachments?: string[];
|
||||
}
|
||||
|
||||
export type ProposalFinishKind = "done" | "cancelled";
|
||||
@@ -325,6 +332,17 @@ function validateProposal(proposal: Proposal, id: string): void {
|
||||
|| typeof pending.receivedAt !== "number") {
|
||||
invalid("pending");
|
||||
}
|
||||
if (pending.attachments !== undefined
|
||||
&& (!Array.isArray(pending.attachments) || pending.attachments.length > 3
|
||||
|| pending.attachments.some((attachment) => !isRecord(attachment)
|
||||
|| typeof attachment.path !== "string" || attachment.path.length === 0
|
||||
|| (attachment.mimeType !== undefined && typeof attachment.mimeType !== "string")))) {
|
||||
invalid("pending.attachments");
|
||||
}
|
||||
if (pending.droppedAttachments !== undefined
|
||||
&& (!Array.isArray(pending.droppedAttachments) || pending.droppedAttachments.some((entry) => typeof entry !== "string"))) {
|
||||
invalid("pending.droppedAttachments");
|
||||
}
|
||||
}
|
||||
if (proposal.finishKind !== undefined && !FINISH_KINDS.includes(proposal.finishKind)) invalid("finishKind");
|
||||
if (proposal.finishNote !== undefined && typeof proposal.finishNote !== "string") invalid("finishNote");
|
||||
|
||||
@@ -28,11 +28,25 @@ export interface IncomingMessage {
|
||||
raw?: unknown;
|
||||
}
|
||||
|
||||
export interface OutgoingImage {
|
||||
mimeType: string;
|
||||
data: string; // base64
|
||||
filename?: string;
|
||||
}
|
||||
|
||||
// A reference to an image file on disk (inside bot.workspace); resolved to OutgoingImage at send time.
|
||||
export interface OutgoingImageRef {
|
||||
path: string;
|
||||
mimeType?: string;
|
||||
filename?: string;
|
||||
}
|
||||
|
||||
export interface OutgoingMessage {
|
||||
target: MessageTarget;
|
||||
text: string;
|
||||
replyTo?: string;
|
||||
replySequence?: number;
|
||||
images?: OutgoingImage[];
|
||||
}
|
||||
|
||||
export interface WebhookRequestContext {
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import type { OutgoingImage, OutgoingImageRef } from "./types.js";
|
||||
|
||||
// Outbound images (Worker screenshots etc.) must live inside bot.workspace and be png or jpg.
|
||||
export const MAX_OUTBOUND_IMAGE_BYTES = 10 * 1024 * 1024;
|
||||
export const MAX_OUTBOUND_IMAGES = 3;
|
||||
|
||||
const PNG_MAGIC = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]);
|
||||
const JPG_MAGIC = Buffer.from([0xff, 0xd8, 0xff]);
|
||||
|
||||
export interface ValidatedWorkspaceImage {
|
||||
path: string; // resolved real path inside the workspace
|
||||
mimeType: "image/png" | "image/jpeg";
|
||||
filename: string;
|
||||
}
|
||||
|
||||
// Resolves a worker/assistant-reported image path against the workspace and validates
|
||||
// containment, size, and magic bytes. Returns undefined (never throws) when invalid.
|
||||
export async function validateWorkspaceImage(workspace: string, refPath: string): Promise<ValidatedWorkspaceImage | undefined> {
|
||||
if (typeof refPath !== "string" || refPath.length === 0) return undefined;
|
||||
try {
|
||||
const root = fs.realpathSync.native(workspace);
|
||||
const candidate = path.resolve(root, refPath);
|
||||
const resolved = await fs.promises.realpath(candidate);
|
||||
const relative = path.relative(root, resolved);
|
||||
if (relative === "" || relative.startsWith("..") || path.isAbsolute(relative)) return undefined;
|
||||
const stat = await fs.promises.lstat(resolved);
|
||||
if (!stat.isFile() || stat.isSymbolicLink() || stat.size === 0 || stat.size > MAX_OUTBOUND_IMAGE_BYTES) return undefined;
|
||||
const handle = await fs.promises.open(resolved, "r");
|
||||
let header: Buffer;
|
||||
try {
|
||||
header = Buffer.alloc(PNG_MAGIC.length);
|
||||
const { bytesRead } = await handle.read(header, 0, PNG_MAGIC.length, 0);
|
||||
header = header.subarray(0, bytesRead);
|
||||
} finally {
|
||||
await handle.close();
|
||||
}
|
||||
const mimeType = magicMimeType(header);
|
||||
if (!mimeType) return undefined;
|
||||
return { path: resolved, mimeType, filename: path.basename(resolved) };
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
// Re-reads a previously validated image at send time (the worker may have updated it);
|
||||
// re-checks type and size and returns undefined when the file is gone or no longer valid.
|
||||
export async function readOutboundImage(ref: OutgoingImageRef): Promise<OutgoingImage | undefined> {
|
||||
try {
|
||||
const stat = await fs.promises.lstat(ref.path);
|
||||
if (!stat.isFile() || stat.isSymbolicLink() || stat.size === 0 || stat.size > MAX_OUTBOUND_IMAGE_BYTES) return undefined;
|
||||
const data = await fs.promises.readFile(ref.path);
|
||||
const mimeType = magicMimeType(data.subarray(0, PNG_MAGIC.length));
|
||||
if (!mimeType) return undefined;
|
||||
return { mimeType, data: data.toString("base64"), filename: ref.filename || path.basename(ref.path) };
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
function magicMimeType(header: Buffer): "image/png" | "image/jpeg" | undefined {
|
||||
if (header.length >= PNG_MAGIC.length && header.subarray(0, PNG_MAGIC.length).equals(PNG_MAGIC)) return "image/png";
|
||||
if (header.length >= JPG_MAGIC.length && header.subarray(0, JPG_MAGIC.length).equals(JPG_MAGIC)) return "image/jpeg";
|
||||
return undefined;
|
||||
}
|
||||
@@ -20,6 +20,7 @@ const ATTACHMENT_DOWNLOAD_TIMEOUT_MS = 10_000;
|
||||
|
||||
export class QqAdapter implements PlatformAdapter {
|
||||
readonly name = "qq";
|
||||
readonly supportsImages = true;
|
||||
private accessToken?: { token: string; expiresAt: number };
|
||||
|
||||
constructor(
|
||||
@@ -68,7 +69,15 @@ export class QqAdapter implements PlatformAdapter {
|
||||
const userOpenId = stringField(raw.user_openid) || stringField(author.user_openid);
|
||||
const isGroup = Boolean(groupOpenId) || message.target.chatId.startsWith("group:");
|
||||
const targetId = groupOpenId || userOpenId || message.target.chatId.replace(/^group:/, "").replace(/^user:/, "");
|
||||
const path = isGroup ? `/v2/groups/${encodeURIComponent(targetId)}/messages` : `/v2/users/${encodeURIComponent(targetId)}/messages`;
|
||||
const base = isGroup ? `/v2/groups/${encodeURIComponent(targetId)}` : `/v2/users/${encodeURIComponent(targetId)}`;
|
||||
|
||||
// Outbound images: upload via /files (group uploads can only be sent to groups, user
|
||||
// uploads only to C2C — the base path already matches the target), then send msg_type 7.
|
||||
for (const image of message.images || []) {
|
||||
await this.sendImageMessage(token, base, isGroup, image, message);
|
||||
}
|
||||
if (message.images?.length && !message.text) return;
|
||||
|
||||
console.log(`QQ send ${isGroup ? "group" : "user"} message length=${message.text.length}`);
|
||||
|
||||
// Passive replies reference the inbound message; proactive messages must send neither msg_id nor msg_seq.
|
||||
@@ -77,7 +86,7 @@ export class QqAdapter implements PlatformAdapter {
|
||||
body.msg_id = message.replyTo;
|
||||
body.msg_seq = message.replySequence ?? 1;
|
||||
}
|
||||
const response = await fetch(`https://api.sgroup.qq.com${path}`, {
|
||||
const response = await fetch(`https://api.sgroup.qq.com${base}/messages`, {
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `QQBot ${token}`,
|
||||
@@ -91,6 +100,37 @@ export class QqAdapter implements PlatformAdapter {
|
||||
if (data?.code && data.code !== 0) throw new Error(`QQ send failed:${qqErrorDetail(data)}`);
|
||||
}
|
||||
|
||||
private async sendImageMessage(token: string, base: string, isGroup: boolean, image: { mimeType: string; data: string; filename?: string }, message: OutgoingMessage): Promise<void> {
|
||||
console.log(`QQ send ${isGroup ? "group" : "user"} image type=${image.mimeType} bytes=${Math.floor(image.data.length * 3 / 4)}`);
|
||||
const headers = {
|
||||
Authorization: `QQBot ${token}`,
|
||||
"Content-Type": "application/json; charset=utf-8"
|
||||
};
|
||||
const upload = await fetch(`https://api.sgroup.qq.com${base}/files`, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body: JSON.stringify({ file_type: 1, file_data: image.data, srv_send_msg: false })
|
||||
});
|
||||
const uploaded = await upload.json().catch(() => undefined) as { file_info?: string; code?: number; message?: string } | undefined;
|
||||
if (!upload.ok) throw new Error(`QQ image upload failed: HTTP ${upload.status}${qqErrorDetail(uploaded)}`);
|
||||
if (uploaded?.code && uploaded.code !== 0) throw new Error(`QQ image upload failed:${qqErrorDetail(uploaded)}`);
|
||||
if (!uploaded?.file_info) throw new Error("QQ image upload failed: missing file_info");
|
||||
|
||||
const body: Record<string, unknown> = { msg_type: 7, media: { file_info: uploaded.file_info }, content: "" };
|
||||
if (message.replyTo) {
|
||||
body.msg_id = message.replyTo;
|
||||
body.msg_seq = message.replySequence ?? 1;
|
||||
}
|
||||
const response = await fetch(`https://api.sgroup.qq.com${base}/messages`, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body: JSON.stringify(body)
|
||||
});
|
||||
const data = await response.json().catch(() => undefined) as QqSendMessageResponse | undefined;
|
||||
if (!response.ok) throw new Error(`QQ image send failed: HTTP ${response.status}${qqErrorDetail(data)}`);
|
||||
if (data?.code && data.code !== 0) throw new Error(`QQ image send failed:${qqErrorDetail(data)}`);
|
||||
}
|
||||
|
||||
private handleValidation(data: QqWebhookEventData | undefined): WebhookResponse {
|
||||
const secret = this.callbackSecret();
|
||||
if (!secret) return jsonResponse({ ok: false, error: "QQ botSecret is required for callback validation" }, 400);
|
||||
|
||||
@@ -2,7 +2,7 @@ import crypto from "node:crypto";
|
||||
import type { AppConfig, BotConfig } from "../config.js";
|
||||
import { SkillLoader, type LoadedSkill } from "./skill-loader.js";
|
||||
|
||||
const BOOTSTRAP_SCHEMA_VERSION = 8;
|
||||
const BOOTSTRAP_SCHEMA_VERSION = 9;
|
||||
|
||||
export interface ResolvedBot extends BotConfig {
|
||||
loadedSkills: LoadedSkill[];
|
||||
@@ -46,8 +46,9 @@ function buildAssistantBootstrap(bot: BotConfig): string {
|
||||
"You are the user-facing Assistant. You have no tools and cannot inspect files, run commands, call Skills or MCP. You chat with the user, propose work, and explain worker feedback. Never claim to have executed anything yourself.",
|
||||
"Speaking style: talk like a reliable colleague, not a console. Lead with the conclusion, then the reason, then the next step. Avoid protocol jargon and field names; never expose internal words like scheduler, pending, envelope, or action types to the user. Default to 2-4 sentences. Do not repeat proposal IDs unless the user asks. If you are unsure, say so plainly. When something is blocked, always give the user an actionable next step.",
|
||||
"Every reply must end with exactly one hidden action envelope: <GORI_ASSISTANT_ACTION_V2>{\"reply\":\"...\",\"actions\":[...]}</GORI_ASSISTANT_ACTION_V2>. Put the user-facing text in the JSON \"reply\" field, not outside the envelope.",
|
||||
"Supported actions: {\"type\":\"create_proposal\",\"title\":\"...\",\"goal\":\"...\",\"steps\":[\"...\"]}; {\"type\":\"confirm\",\"id\":\"optional\"}; {\"type\":\"adjust_proposal\",\"id\":\"optional\",\"title\":\"optional\",\"goal\":\"optional\",\"steps\":\"optional\"}; {\"type\":\"follow_up\",\"id\":\"optional\",\"instruction\":\"...\"}; {\"type\":\"finish\",\"id\":\"optional\",\"note\":\"optional\"}; {\"type\":\"start_next\"}; {\"type\":\"cancel\",\"id\":\"optional\"}; {\"type\":\"stop\"}. Use an empty actions array when no state change is needed.",
|
||||
"Supported actions: {\"type\":\"create_proposal\",\"title\":\"...\",\"goal\":\"...\",\"steps\":[\"...\"]}; {\"type\":\"confirm\",\"id\":\"optional\"}; {\"type\":\"adjust_proposal\",\"id\":\"optional\",\"title\":\"optional\",\"goal\":\"optional\",\"steps\":\"optional\"}; {\"type\":\"follow_up\",\"id\":\"optional\",\"instruction\":\"...\"}; {\"type\":\"finish\",\"id\":\"optional\",\"note\":\"optional\"}; {\"type\":\"send_image\",\"path\":\"...\"}; {\"type\":\"start_next\"}; {\"type\":\"cancel\",\"id\":\"optional\"}; {\"type\":\"stop\"}. Use an empty actions array when no state change is needed.",
|
||||
"A proposal only starts after the user confirms it. When a worker finishes a turn, the proposal becomes pending: it waits for the user's decision with a summary (and maybe a question). The worker never reports success or failure; treat every result as information for the user. \"finish\" closes a pending proposal as done; \"follow_up\" sends the user's new instruction to the same proposal and resumes its worker; \"cancel\" drops a proposed, queued, or pending proposal; \"stop\" aborts the running worker and leaves the proposal pending; \"start_next\" starts the oldest confirmed queued proposal. \"adjust_proposal\" edits a proposal that has not started yet. Never invent other actions or statuses.",
|
||||
"When a worker's summary says it saved image files inside the workspace (for example under .gori-outbox/) and the user asks to see one, emit \"send_image\" with that exact path. Only send paths a worker actually reported; never invent paths.",
|
||||
"Whenever the [Pending proposals] section in a prompt lists one of the user's proposals, your reply MUST acknowledge it: remind the user what is waiting and that they can say finish to close it or just keep talking to continue it.",
|
||||
"Never claim an action has already taken effect. The runtime executes your actions after your reply and appends a correction to your message when something could not be done (for example when start_next is blocked by another proposal). Treat that correction as the truth and use the [Proposal states], [Pending proposals], [Scheduler state] and [Worker state] sections in each prompt as the only reliable state.",
|
||||
"In a group chat each proposal is owned by the user who requested it: only that user can confirm, adjust, follow up, finish, stop, or cancel it. If another group member asks you to operate a proposal they did not create, explain that only the proposal's initiator can do that and emit no action.",
|
||||
@@ -63,7 +64,7 @@ function buildWorkerBootstrap(bot: BotConfig, skills: LoadedSkill[]): string {
|
||||
bot.persona ? `Persona:\n${bot.persona}` : "Persona: general coding assistant",
|
||||
`Permission policy enforced by the ACP client: ${JSON.stringify(bot.permissions)}`,
|
||||
"You are the Worker. You execute exactly one confirmed proposal in the configured workspace and never talk to the user directly.",
|
||||
"For every turn, end your response with exactly one hidden result envelope: <GORI_WORKER_RESULT_V2>{\"status\":\"PENDING\",\"summary\":\"...\"}</GORI_WORKER_RESULT_V2>. PENDING is the only status: it hands the result back to the user. Always include a short user-readable \"summary\" of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing. Set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. Never emit any other status or text after the envelope.",
|
||||
"For every turn, end your response with exactly one hidden result envelope: <GORI_WORKER_RESULT_V2>{\"status\":\"PENDING\",\"summary\":\"...\"}</GORI_WORKER_RESULT_V2>. PENDING is the only status: it hands the result back to the user. Always include a short user-readable \"summary\" of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing. Set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. When you produced image files the user should see (png/jpg only), save them inside the workspace (prefer .gori-outbox/) and report up to 3 as \"attachments\": [{\"path\":\"relative/or/absolute/path\",\"mimeType\":\"optional\"}]; never report paths outside the workspace such as /tmp. Never emit any other status or text after the envelope.",
|
||||
"Keep every command and tool process attached to this ACP worker. Never daemonize, call setsid, use nohup, create a detached process, or leave a background process running after the turn.",
|
||||
...skills.map((skill) => `Skill ${skill.id} (${skill.file}):\n${skill.content}`)
|
||||
].join("\n\n");
|
||||
|
||||
+1
-1
@@ -47,7 +47,7 @@ export async function createGatewayRuntime(config: AppConfig): Promise<GatewayRu
|
||||
await assistantManager.initialize();
|
||||
const manager = assistantManager;
|
||||
const gateway = new Gateway(config.gateway.policy, manager);
|
||||
manager.setEventSink((chatKey, text) => gateway.sendEvent(chatKey, text));
|
||||
manager.setEventSink((chatKey, text, images) => gateway.sendEvent(chatKey, text, images));
|
||||
const platformAdapter = createPlatformAdapter(config, gateway);
|
||||
const qqGatewayClient = config.gateway.platform.type === "qq" && config.gateway.platform.connectionMode === "websocket"
|
||||
? new QqGatewayClient(config.gateway.platform, platformAdapter as QqAdapter) : undefined;
|
||||
|
||||
Reference in New Issue
Block a user