feat(agent): persist message metadata and previews

This commit is contained in:
yu
2026-08-06 19:11:41 +08:00
parent 89a4e192e5
commit 59424995d7
11 changed files with 520 additions and 141 deletions
+1
View File
@@ -84,4 +84,5 @@
- 当前 AI API Key 存在浏览器本地,并由前端直接请求 OpenAI 兼容接口;涉及安全说明时要写清楚。 - 当前 AI API Key 存在浏览器本地,并由前端直接请求 OpenAI 兼容接口;涉及安全说明时要写清楚。
- Docker 静态资源路径目前仍是待办项,文档中不要过度承诺生产部署已经完全验证。 - Docker 静态资源路径目前仍是待办项,文档中不要过度承诺生产部署已经完全验证。
- Agent 对话消息必须同时按 `threadId`、`turnId` 和 `itemId` 归属;实时事件只用于补充未物化的 turn,历史快照成为权威后不得重复合并同一条消息。 - Agent 对话消息必须同时按 `threadId`、`turnId` 和 `itemId` 归属;实时事件只用于补充未物化的 turn,历史快照成为权威后不得重复合并同一条消息。
- Agent 通信协议版本与消息存储版本必须独立管理;消息存储格式升级时必须先备份再迁移,遇到未知版本、损坏清单或冲突备份时拒绝覆盖原文件,不得按记录数量或文件大小静默裁剪历史元数据。
- 本地启动或浏览器验收时不要关闭用户已经打开的浏览器窗口或标签页;需要自动化验证时使用独立测试页面,避免打断用户当前页面和对话状态。 - 本地启动或浏览器验收时不要关闭用户已经打开的浏览器窗口或标签页;需要自动化验证时使用独立测试页面,避免打断用户当前页面和对话状态。
+1 -1
View File
@@ -15,7 +15,7 @@
"scripts": { "scripts": {
"dev": "tsx src/index.ts", "dev": "tsx src/index.ts",
"debug": "tsx src/index.ts --debug", "debug": "tsx src/index.ts --debug",
"test": "tsx --test src/canvas/session.test.ts src/agent/codex-client.test.ts src/agent/codex-history.test.ts src/skills/store.test.ts", "test": "tsx --test src/canvas/session.test.ts src/agent/codex-client.test.ts src/agent/codex-history.test.ts src/agent/message-metadata.test.ts src/skills/store.test.ts",
"build": "tsc -p tsconfig.json", "build": "tsc -p tsconfig.json",
"start": "node dist/index.js", "start": "node dist/index.js",
"prepack": "npm run build" "prepack": "npm run build"
+15 -2
View File
@@ -9,6 +9,7 @@ import { errorMessage, field, type JsonRecord } from "../utils/value.js";
import { CodexAppClient, CodexReportedError } from "./codex-client.js"; import { CodexAppClient, CodexReportedError } from "./codex-client.js";
import { codexEventHistory } from "./codex-event-history.js"; import { codexEventHistory } from "./codex-event-history.js";
import { settledTurnIds, summarizeCodexThread, threadMessages } from "./codex-history.js"; import { settledTurnIds, summarizeCodexThread, threadMessages } from "./codex-history.js";
import { messageMetadataStore } from "./message-metadata.js";
import type { CodexReasoningEffort, CodexSkillMetadata, CodexSkillSelector, CodexSkillsListEntry } from "./codex-protocol.js"; import type { CodexReasoningEffort, CodexSkillMetadata, CodexSkillSelector, CodexSkillsListEntry } from "./codex-protocol.js";
import type { AgentAttachment, AgentEmit, AgentPermissionMode } from "./types.js"; import type { AgentAttachment, AgentEmit, AgentPermissionMode } from "./types.js";
@@ -97,7 +98,8 @@ export async function resumeCodexThread(emit: AgentEmit, threadId: string, cwd?:
const thread = await resumeLoadedThread(app, threadId, cwd, permissionMode, true, preheat); const thread = await resumeLoadedThread(app, threadId, cwd, permissionMode, true, preheat);
const history = await loadCodexHistory(emit, threadId, cwd); const history = await loadCodexHistory(emit, threadId, cwd);
const supplementalItems = await codexEventHistory.readThread(threadId); const supplementalItems = await codexEventHistory.readThread(threadId);
return { thread, messages: threadMessages(history.thread, app.planUpdates(threadId), supplementalItems), settledTurnIds: settledTurnIds(history.thread, supplementalItems), historyReady: history.historyReady }; const messages = await mergeMessageMetadata(threadId, threadMessages(history.thread, app.planUpdates(threadId), supplementalItems));
return { thread, messages, settledTurnIds: settledTurnIds(history.thread, supplementalItems), historyReady: history.historyReady };
} }
/** 查询当前工作空间中的 Codex 线程。 */ /** 查询当前工作空间中的 Codex 线程。 */
@@ -150,7 +152,8 @@ export async function readCodexThread(emit: AgentEmit, threadId: string, cwd?: s
const app = await getCodexApp(emit); const app = await getCodexApp(emit);
const history = await loadCodexHistory(emit, threadId, cwd); const history = await loadCodexHistory(emit, threadId, cwd);
const supplementalItems = await codexEventHistory.readThread(threadId); const supplementalItems = await codexEventHistory.readThread(threadId);
return { thread: summarizeCodexThread(history.thread), messages: threadMessages(history.thread, app.planUpdates(threadId), supplementalItems), settledTurnIds: settledTurnIds(history.thread, supplementalItems), historyReady: history.historyReady }; const messages = await mergeMessageMetadata(threadId, threadMessages(history.thread, app.planUpdates(threadId), supplementalItems));
return { thread: summarizeCodexThread(history.thread), messages, settledTurnIds: settledTurnIds(history.thread, supplementalItems), historyReady: history.historyReady };
} }
/** 归档指定 Codex 线程。 */ /** 归档指定 Codex 线程。 */
@@ -165,9 +168,19 @@ export async function archiveCodexThread(emit: AgentEmit, threadId: string, cwd?
await app.archiveThread(threadId); await app.archiveThread(threadId);
app.clearPlanUpdates(threadId); app.clearPlanUpdates(threadId);
await codexEventHistory.removeThread(threadId); await codexEventHistory.removeThread(threadId);
await messageMetadataStore.removeThread(threadId).catch((error) => logger.warn("Failed to remove archived thread message metadata", { threadId, error }));
if (loadedThreadId === threadId) loadedThreadId = ""; if (loadedThreadId === threadId) loadedThreadId = "";
} }
async function mergeMessageMetadata<T extends { role: string; threadId: string; turnId: string }>(threadId: string, messages: T[]) {
try {
return await messageMetadataStore.mergeThread(threadId, messages);
} catch (error) {
logger.warn("Failed to read thread message metadata", { threadId, error });
return messages;
}
}
/** 判断线程异常是否允许自动新建线程后重试。 */ /** 判断线程异常是否允许自动新建线程后重试。 */
export function isRecoverableThreadError(error: unknown) { export function isRecoverableThreadError(error: unknown) {
return /thread not loaded|no rollout found/i.test(errorMessage(error)); return /thread not loaded|no rollout found/i.test(errorMessage(error));
@@ -0,0 +1,104 @@
import assert from "node:assert/strict";
import fs from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import test, { type TestContext } from "node:test";
import { MessageMetadataStore } from "./message-metadata.js";
test("message metadata survives Agent restarts", async (context) => {
const fixture = await createFixture(context);
await fixture.store.recordPending("message-1", sampleMetadata());
await fixture.store.bindThread("message-1", "thread-1");
await fixture.store.bindTurn("message-1", "thread-1", "turn-1");
const [message] = await fixture.reopen().mergeThread("thread-1", [{ role: "user", threadId: "thread-1", turnId: "turn-1", text: "Generate product images" }]);
assert.equal(message.clientMessageId, "message-1");
assert.equal(message.attachments?.[0].name, "image.png");
assert.equal(message.canvasReferences?.[0].nodeId, "node-1");
assert.equal(message.skill?.name, "product-grid");
});
test("message preview assets survive restarts and are deleted with their thread", async (context) => {
const fixture = await createFixture(context);
const metadata = await fixture.store.recordPending("message-1", sampleMetadata());
await fixture.store.bindTurn("message-1", "thread-1", "turn-1");
const match = metadata?.attachments?.[0].url.match(/^agent-asset:([a-f0-9]{64})\/([a-f0-9]{64}\.png)$/);
assert.ok(match);
const asset = await fixture.reopen().readAsset(match[1], match[2]);
assert.equal(asset?.contentType, "image/png");
assert.equal(asset?.data.toString(), "a");
await fixture.store.removeThread("thread-1");
assert.equal(await fixture.reopen().readAsset(match[1], match[2]), undefined);
});
test("deleting a thread only removes its metadata", async (context) => {
const fixture = await createFixture(context);
for (const number of [1, 2]) {
await fixture.store.recordPending(`message-${number}`, { skill: { name: `skill-${number}`, path: `D:\\skills\\skill-${number}\\SKILL.md` } });
await fixture.store.bindTurn(`message-${number}`, `thread-${number}`, `turn-${number}`);
}
await fixture.store.removeThread("thread-1");
const first = await fixture.reopen().mergeThread("thread-1", [{ role: "user", threadId: "thread-1", turnId: "turn-1" }]);
const second = await fixture.reopen().mergeThread("thread-2", [{ role: "user", threadId: "thread-2", turnId: "turn-2" }]);
assert.equal(first[0].skill, undefined);
assert.equal(second[0].skill?.name, "skill-2");
});
test("history metadata is matched by thread and turn instead of client message id alone", async (context) => {
const fixture = await createFixture(context);
await fixture.store.recordPending("message-1", { skill: { name: "product-grid", path: "D:\\skills\\product-grid\\SKILL.md" } });
await fixture.store.bindTurn("message-1", "thread-1", "turn-1");
const [message] = await fixture.store.mergeThread("thread-1", [{ role: "user", threadId: "thread-1", turnId: "turn-2", clientMessageId: "message-1" }]);
assert.equal(message.skill, undefined);
});
test("unknown storage versions are never overwritten", async (context) => {
const fixture = await createFixture(context, false);
await fs.mkdir(fixture.storeDirectory, { recursive: true });
const manifestFile = path.join(fixture.storeDirectory, "manifest.json");
await fs.writeFile(manifestFile, '{"version":99}');
await assert.rejects(() => fixture.store.recordPending("message-1", sampleMetadata()), /Unsupported message metadata version/);
assert.equal(await fs.readFile(manifestFile, "utf8"), '{"version":99}');
});
test("storage without a manifest is never overwritten", async (context) => {
const fixture = await createFixture(context, false);
await fs.mkdir(fixture.storeDirectory, { recursive: true });
const existingFile = path.join(fixture.storeDirectory, "existing.json");
await fs.writeFile(existingFile, "existing data");
await assert.rejects(() => fixture.store.mergeThread("thread-1", []), /missing manifest/);
assert.equal(await fs.readFile(existingFile, "utf8"), "existing data");
});
test("oversized image previews are rejected instead of silently dropped", async (context) => {
const fixture = await createFixture(context);
await assert.rejects(
() => fixture.store.recordPending("message-1", { attachments: [{ id: "image-1", name: "large.png", url: `data:image/png;base64,${"a".repeat(500_001)}` }] }),
/attachments metadata is invalid/,
);
});
async function createFixture(context: TestContext, initialize = true) {
const root = await fs.mkdtemp(path.join(os.tmpdir(), "canvas-agent-message-metadata-"));
context.after(() => fs.rm(root, { recursive: true, force: true }));
const storeDirectory = path.join(root, "message-metadata");
const createStore = () => new MessageMetadataStore(storeDirectory);
const fixture = { root, storeDirectory, store: createStore(), reopen: createStore };
if (initialize) await fixture.store.mergeThread("empty", []);
return fixture;
}
function sampleMetadata() {
return {
attachments: [{ id: "image-1", name: "image.png", url: "data:image/png;base64,YQ==" }],
canvasReferences: [{ nodeId: "node-1", label: "Image 1", title: "Product image", kind: "image", previewUrl: "data:image/png;base64,YQ==" }],
skill: { name: "product-grid", path: "D:\\skills\\product-grid\\SKILL.md", displayName: "Product grid" },
};
}
+369
View File
@@ -0,0 +1,369 @@
import crypto from "node:crypto";
import fs from "node:fs/promises";
import path from "node:path";
import { CONFIG_DIR } from "../config.js";
import type { AgentAttachmentDisplay, AgentCanvasReference, AgentMessageMetadata, AgentSkillReference } from "./types.js";
type MessageMetadataRecord = { version: 1; clientMessageId: string; threadId?: string; turnId?: string; createdAt: number; metadata: AgentMessageMetadata };
type MetadataMessage = { role: string; threadId: string; turnId: string; clientMessageId?: string };
const STORAGE_VERSION = 1;
const MAX_PREVIEW_LENGTH = 500_000;
const MAX_TEXT_LENGTH = 20_000;
const MANIFEST_FILE = "manifest.json";
const MESSAGE_ASSET_PREFIX = "agent-asset:";
const MESSAGE_METADATA_DIRECTORY = path.join(CONFIG_DIR, "message-metadata");
export class MessageMetadataStore {
private queue: Promise<void> = Promise.resolve();
private ready?: Promise<void>;
constructor(private directory = MESSAGE_METADATA_DIRECTORY) {}
recordPending(clientMessageId: string, value: unknown) {
return this.run(async () => {
const metadata = normalizeMetadata(value, false);
if (!clientMessageId || !metadata) return undefined;
await this.ensureReady();
const storedMetadata = await persistMetadataPreviews(this.directory, clientMessageId, metadata);
const record: MessageMetadataRecord = { version: STORAGE_VERSION, clientMessageId, createdAt: Date.now(), metadata: storedMetadata };
await writeJson(this.pendingFile(clientMessageId), record);
return structuredClone(storedMetadata);
});
}
bindThread(clientMessageId: string, threadId: string) {
return this.run(async () => {
if (!clientMessageId || !threadId) return;
await this.ensureReady();
const pendingFile = this.pendingFile(clientMessageId);
const targetFile = this.threadFile(threadId, clientMessageId);
const record = await readRecord(targetFile) || await readRecord(pendingFile);
if (!record) return;
await writeJson(targetFile, { ...record, threadId, turnId: record.threadId === threadId ? record.turnId : undefined });
await removeFile(pendingFile);
});
}
bindTurn(clientMessageId: string, threadId: string, turnId: string) {
return this.run(async () => {
if (!clientMessageId || !threadId || !turnId) return;
await this.ensureReady();
const pendingFile = this.pendingFile(clientMessageId);
const targetFile = this.threadFile(threadId, clientMessageId);
const record = await readRecord(targetFile) || await readRecord(pendingFile);
if (!record) return;
await writeJson(targetFile, { ...record, threadId, turnId });
await removeFile(pendingFile);
});
}
mergeThread<T extends MetadataMessage>(threadId: string, messages: T[]) {
return this.run<Array<T & AgentMessageMetadata>>(async () => {
await this.ensureReady();
const records = (await readRecords(this.threadDirectory(threadId))).filter((item) => item.threadId === threadId && item.turnId);
const byTurnId = new Map(records.map((item) => [item.turnId!, item]));
return messages.map((message) => {
if (message.role !== "user" || message.threadId !== threadId) return message;
const record = byTurnId.get(message.turnId);
return record ? { ...message, clientMessageId: record.clientMessageId, ...structuredClone(record.metadata) } : message;
});
});
}
remove(clientMessageId: string, threadId = "") {
return this.run(async () => {
if (!clientMessageId) return;
await this.ensureReady();
await Promise.all([
removeFile(this.pendingFile(clientMessageId)),
...(threadId ? [removeFile(this.threadFile(threadId, clientMessageId))] : []),
fs.rm(this.assetDirectory(clientMessageId), { recursive: true, force: true }),
]);
});
}
removeThread(threadId: string) {
return this.run(async () => {
if (!threadId) return;
await this.ensureReady();
const records = await readRecords(this.threadDirectory(threadId));
await Promise.all([
fs.rm(this.threadDirectory(threadId), { recursive: true, force: true }),
...records.map((record) => fs.rm(this.assetDirectory(record.clientMessageId), { recursive: true, force: true })),
]);
});
}
readAsset(messageKey: string, assetFile: string) {
return this.run(async () => {
await this.ensureReady();
if (!/^[a-f0-9]{64}$/.test(messageKey) || !/^[a-f0-9]{64}\.(?:gif|jpe?g|png|webp)$/.test(assetFile)) return undefined;
try {
return { data: await fs.readFile(path.join(this.directory, "assets", messageKey, assetFile)), contentType: assetContentType(assetFile) };
} catch (error) {
if (isMissing(error)) return undefined;
throw error;
}
});
}
private run<T>(task: () => Promise<T>) {
const result = this.queue.then(task, task);
this.queue = result.then(() => undefined, () => undefined);
return result;
}
private ensureReady() {
this.ready ||= this.initialize();
return this.ready;
}
private async initialize() {
const manifestFile = path.join(this.directory, MANIFEST_FILE);
try {
const manifest = await readJson(manifestFile) as { version?: unknown };
if (manifest.version !== STORAGE_VERSION) throw unsupportedVersion(manifest.version);
return;
} catch (error) {
if (!isMissing(error)) throw error;
}
if (await exists(this.directory)) throw new Error(`Message metadata store is missing ${MANIFEST_FILE}; refusing to overwrite it.`);
await createStorage(this.directory);
}
private pendingFile(clientMessageId: string) {
return path.join(this.directory, "pending", `${storageKey(clientMessageId)}.json`);
}
private threadDirectory(threadId: string) {
return path.join(this.directory, "threads", storageKey(threadId));
}
private threadFile(threadId: string, clientMessageId: string) {
return path.join(this.threadDirectory(threadId), `${storageKey(clientMessageId)}.json`);
}
private assetDirectory(clientMessageId: string) {
return path.join(this.directory, "assets", storageKey(clientMessageId));
}
}
export const messageMetadataStore = new MessageMetadataStore();
async function createStorage(directory: string) {
const temporaryDirectory = `${directory}.${process.pid}.${Date.now()}.tmp`;
try {
await fs.mkdir(path.join(temporaryDirectory, "pending"), { recursive: true });
await fs.mkdir(path.join(temporaryDirectory, "threads"), { recursive: true });
await fs.mkdir(path.join(temporaryDirectory, "assets"), { recursive: true });
await writeJson(path.join(temporaryDirectory, MANIFEST_FILE), { version: STORAGE_VERSION });
await fs.rename(temporaryDirectory, directory);
} finally {
await fs.rm(temporaryDirectory, { recursive: true, force: true });
}
}
async function readRecords(directory: string) {
let files: string[];
try {
files = (await fs.readdir(directory, { withFileTypes: true })).filter((item) => item.isFile() && item.name.endsWith(".json")).map((item) => path.join(directory, item.name));
} catch (error) {
if (isMissing(error)) return [];
throw error;
}
return (await Promise.all(files.map(readRecord))).filter((item): item is MessageMetadataRecord => Boolean(item));
}
async function readRecord(file: string): Promise<MessageMetadataRecord | undefined> {
try {
const value = await readJson(file) as Partial<MessageMetadataRecord>;
if (value.version !== STORAGE_VERSION) throw unsupportedVersion(value.version);
const clientMessageId = text(value.clientMessageId, 200);
const metadata = normalizeMetadata(value.metadata, true);
if (!clientMessageId || !metadata) throw new Error(`Message metadata record is invalid: ${file}`);
const threadId = text(value.threadId, 200);
const turnId = text(value.turnId, 200);
return { version: STORAGE_VERSION, clientMessageId, createdAt: Number(value.createdAt) || 0, metadata, ...(threadId ? { threadId } : {}), ...(turnId ? { turnId } : {}) };
} catch (error) {
if (isMissing(error)) return undefined;
throw error;
}
}
function normalizeMetadata(value: unknown, allowAsset: boolean): AgentMessageMetadata | undefined {
if (value === undefined || value === null) return undefined;
if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("Message metadata must be an object.");
const source = value as Record<string, unknown>;
const attachments = normalizeList(source.attachments, 6, (item) => normalizeAttachment(item, allowAsset), "attachments");
const canvasReferences = normalizeList(source.canvasReferences, 50, (item) => normalizeCanvasReference(item, allowAsset), "canvas references");
const skill = normalizeSkill(source.skill);
if (source.skill !== undefined && !skill) throw new Error("Message skill metadata is invalid.");
if (!attachments.length && !canvasReferences.length && !skill) return undefined;
return { ...(attachments.length ? { attachments } : {}), ...(canvasReferences.length ? { canvasReferences } : {}), ...(skill ? { skill } : {}) };
}
function normalizeList<T>(value: unknown, limit: number, normalize: (item: unknown) => T | undefined, label: string) {
if (value === undefined) return [];
if (!Array.isArray(value)) throw new Error(`Message ${label} metadata must be an array.`);
if (value.length > limit) throw new Error(`Message ${label} metadata exceeds the limit of ${limit}.`);
const normalized = value.map(normalize).filter((item): item is T => Boolean(item));
if (normalized.length !== value.length) throw new Error(`Message ${label} metadata is invalid.`);
return normalized;
}
function normalizeAttachment(value: unknown, allowAsset: boolean): AgentAttachmentDisplay | undefined {
if (!value || typeof value !== "object" || Array.isArray(value)) return undefined;
const item = value as Record<string, unknown>;
const id = text(item.id, 200);
const name = text(item.name, 500);
const url = imagePreview(item.url, allowAsset);
if (!id || !name || !url) return undefined;
return { id, name, url, ...numberFields(item) };
}
function normalizeCanvasReference(value: unknown, allowAsset: boolean): AgentCanvasReference | undefined {
if (!value || typeof value !== "object" || Array.isArray(value)) return undefined;
const item = value as Record<string, unknown>;
const nodeId = text(item.nodeId, 200);
const label = text(item.label, 200);
const title = text(item.title, 500);
const kind = String(item.kind || "");
if (!nodeId || !label || !title || !["image", "video", "audio", "text"].includes(kind)) return undefined;
const previewUrl = kind === "image" ? imagePreview(item.previewUrl, allowAsset) : remotePreview(item.previewUrl);
const referenceText = text(item.text, MAX_TEXT_LENGTH);
return { nodeId, label, title, kind: kind as AgentCanvasReference["kind"], ...(previewUrl ? { previewUrl } : {}), ...(referenceText ? { text: referenceText } : {}) };
}
function normalizeSkill(value: unknown): AgentSkillReference | undefined {
if (value === undefined) return undefined;
if (!value || typeof value !== "object" || Array.isArray(value)) return undefined;
const item = value as Record<string, unknown>;
const name = text(item.name, 200);
const skillPath = text(item.path, 2000);
const displayName = text(item.displayName, 200);
return name && skillPath ? { name, path: skillPath, ...(displayName ? { displayName } : {}) } : undefined;
}
function numberFields(item: Record<string, unknown>) {
const fields: Partial<Pick<AgentAttachmentDisplay, "type" | "size" | "width" | "height">> = {};
const type = text(item.type, 200);
if (type) fields.type = type;
for (const key of ["size", "width", "height"] as const) {
const value = Number(item[key]);
if (Number.isFinite(value) && value >= 0) fields[key] = value;
}
return fields;
}
function imagePreview(value: unknown, allowAsset = false) {
if (typeof value !== "string") return "";
const url = value.trim();
if (allowAsset && new RegExp(`^${MESSAGE_ASSET_PREFIX}[a-f0-9]{64}/[a-f0-9]{64}\\.(?:gif|jpe?g|png|webp)$`).test(url)) return url;
if (/^https?:\/\//i.test(url)) return url.length <= 5000 ? url : "";
return url.length <= MAX_PREVIEW_LENGTH && /^data:image\/[a-z0-9.+-]+;base64,/i.test(url) ? url : "";
}
function remotePreview(value: unknown) {
if (typeof value !== "string") return "";
const url = value.trim();
return url.length <= 5000 && /^https?:\/\//i.test(url) ? url : "";
}
function text(value: unknown, limit: number) {
return typeof value === "string" ? value.trim().slice(0, limit) : "";
}
function storageKey(value: string) {
return crypto.createHash("sha256").update(value).digest("hex");
}
async function persistMetadataPreviews(directory: string, clientMessageId: string, metadata: AgentMessageMetadata) {
const attachments = await Promise.all((metadata.attachments || []).map(async (item) => ({ ...item, url: await persistImagePreview(directory, clientMessageId, `attachment:${item.id}`, item.url) })));
const canvasReferences = await Promise.all((metadata.canvasReferences || []).map(async (item) => ({
...item,
...(item.kind === "image" && item.previewUrl ? { previewUrl: await persistImagePreview(directory, clientMessageId, `reference:${item.nodeId}`, item.previewUrl) } : {}),
})));
return {
...(attachments.length ? { attachments } : {}),
...(canvasReferences.length ? { canvasReferences } : {}),
...(metadata.skill ? { skill: metadata.skill } : {}),
};
}
async function persistImagePreview(directory: string, clientMessageId: string, key: string, value: string) {
const match = value.match(/^data:(image\/[a-z0-9.+-]+);base64,(.+)$/i);
if (!match) return value;
const extension = imageExtension(match[1]);
if (!extension) throw new Error(`Unsupported message preview type: ${match[1]}`);
const messageKey = storageKey(clientMessageId);
const assetFile = `${storageKey(key)}.${extension}`;
await writeFile(path.join(directory, "assets", messageKey, assetFile), Buffer.from(match[2], "base64"));
return `${MESSAGE_ASSET_PREFIX}${messageKey}/${assetFile}`;
}
function imageExtension(contentType: string) {
const type = contentType.toLowerCase();
if (type === "image/jpeg" || type === "image/jpg") return "jpg";
if (type === "image/png") return "png";
if (type === "image/webp") return "webp";
if (type === "image/gif") return "gif";
return "";
}
function assetContentType(file: string) {
const extension = path.extname(file).slice(1).toLowerCase();
if (extension === "jpg" || extension === "jpeg") return "image/jpeg";
return `image/${extension}`;
}
async function readJson(file: string) {
try {
return JSON.parse(await fs.readFile(file, "utf8")) as unknown;
} catch (error) {
if (error instanceof SyntaxError) throw new Error(`Message metadata JSON is invalid: ${file}`);
throw error;
}
}
async function writeJson(file: string, value: unknown) {
await writeFile(file, JSON.stringify(value));
}
async function writeFile(file: string, value: string | Buffer) {
await fs.mkdir(path.dirname(file), { recursive: true });
const temporaryFile = `${file}.${process.pid}.${crypto.randomUUID()}.tmp`;
try {
await fs.writeFile(temporaryFile, value);
await fs.rename(temporaryFile, file);
} finally {
await removeFile(temporaryFile);
}
}
async function removeFile(file: string) {
await fs.unlink(file).catch((error: NodeJS.ErrnoException) => {
if (error.code !== "ENOENT") throw error;
});
}
async function exists(file: string) {
try {
await fs.access(file);
return true;
} catch (error) {
if (isMissing(error)) return false;
throw error;
}
}
function unsupportedVersion(version: unknown) {
return new Error(`Unsupported message metadata version: ${String(version ?? "missing")}. Refusing to overwrite existing data.`);
}
function isMissing(error: unknown) {
return (error as NodeJS.ErrnoException)?.code === "ENOENT";
}
+12
View File
@@ -4,5 +4,17 @@ export type AgentEmit = (type: string, payload: unknown) => void;
/** 用户随当前 Agent 消息上传的附件。 */ /** 用户随当前 Agent 消息上传的附件。 */
export type AgentAttachment = { id?: string; name?: string; type?: string; size?: number; width?: number; height?: number; dataUrl?: string }; export type AgentAttachment = { id?: string; name?: string; type?: string; size?: number; width?: number; height?: number; dataUrl?: string };
/** 用户消息在历史记录中展示所需的附件信息。 */
export type AgentAttachmentDisplay = { id: string; name: string; type?: string; size?: number; width?: number; height?: number; url: string };
/** 用户消息引用的画布素材。 */
export type AgentCanvasReference = { nodeId: string; label: string; title: string; kind: "image" | "video" | "audio" | "text"; previewUrl?: string; text?: string };
/** 用户消息调用的 Codex Skill。 */
export type AgentSkillReference = { name: string; path: string; displayName?: string };
/** 独立于 Codex 原生线程历史保存的用户消息展示元数据。 */
export type AgentMessageMetadata = { attachments?: AgentAttachmentDisplay[]; canvasReferences?: AgentCanvasReference[]; skill?: AgentSkillReference };
/** Codex 文件、命令和网络权限模式。 */ /** Codex 文件、命令和网络权限模式。 */
export type AgentPermissionMode = "request" | "automatic" | "full"; export type AgentPermissionMode = "request" | "automatic" | "full";
+1 -1
View File
@@ -221,7 +221,7 @@ test("new clients receive the current Codex state and later updates", (t) => {
t.after(() => client.close()); t.after(() => client.close());
const hello = client.event("hello"); const hello = client.event("hello");
assert.equal(field(hello, "protocolVersion"), 5); assert.equal(field(hello, "protocolVersion"), 6);
assert.deepEqual(field(hello, "workspace"), { activeThreadId: "thread-2" }); assert.deepEqual(field(hello, "workspace"), { activeThreadId: "thread-2" });
assert.deepEqual(field(hello, "conversation"), { revision: 1, conversationId: "thread-2", threadId: "thread-2", status: "ready", mcpStatuses: {} }); assert.deepEqual(field(hello, "conversation"), { revision: 1, conversationId: "thread-2", threadId: "thread-2", status: "ready", mcpStatuses: {} });
assert.deepEqual(field(hello, "codex"), { busy: true, threadId: "thread-2", turnId: "turn-1" }); assert.deepEqual(field(hello, "codex"), { busy: true, threadId: "thread-2", turnId: "turn-1" });
+1 -1
View File
@@ -23,7 +23,7 @@ export type ConversationState = {
error?: string; error?: string;
}; };
type McpInventoryItem = { name: string; authStatus?: string }; type McpInventoryItem = { name: string; authStatus?: string };
export const AGENT_PROTOCOL_VERSION = 5; export const AGENT_PROTOCOL_VERSION = 6;
const SITE_TOOLS = new Set<ToolName>([ const SITE_TOOLS = new Set<ToolName>([
"site_navigate", "site_navigate",
+15 -3
View File
@@ -6,6 +6,7 @@ import express, { type NextFunction, type Request, type Response } from "express
import { runClaudeTurn } from "../agent/claude.js"; import { runClaudeTurn } from "../agent/claude.js";
import { archiveCodexThread, CodexSkillLookupError, configureCodexSkill, generateCodexSkillDraft, interruptCodexTurn, isRecoverableThreadError, listCodexModels, listCodexSkills, listCodexThreads, readCodexThread, resolveCodexApproval, resolveCodexSkill, resumeCodexThread, runCodexTurn, startCodexThread, summarizeCodexThread } from "../agent/codex.js"; import { archiveCodexThread, CodexSkillLookupError, configureCodexSkill, generateCodexSkillDraft, interruptCodexTurn, isRecoverableThreadError, listCodexModels, listCodexSkills, listCodexThreads, readCodexThread, resolveCodexApproval, resolveCodexSkill, resumeCodexThread, runCodexTurn, startCodexThread, summarizeCodexThread } from "../agent/codex.js";
import type { CodexReasoningEffort, CodexSkillSelector } from "../agent/codex-protocol.js"; import type { CodexReasoningEffort, CodexSkillSelector } from "../agent/codex-protocol.js";
import { messageMetadataStore } from "../agent/message-metadata.js";
import type { AgentAttachment, AgentPermissionMode } from "../agent/types.js"; import type { AgentAttachment, AgentPermissionMode } from "../agent/types.js";
import { AGENT_PROTOCOL_VERSION, CanvasSession } from "../canvas/session.js"; import { AGENT_PROTOCOL_VERSION, CanvasSession } from "../canvas/session.js";
import { DEFAULT_PORT, ensureSiteWorkspace, loadConfig, saveConfig, updateSiteWorkspace, type CanvasAgentConfig } from "../config.js"; import { DEFAULT_PORT, ensureSiteWorkspace, loadConfig, saveConfig, updateSiteWorkspace, type CanvasAgentConfig } from "../config.js";
@@ -147,6 +148,12 @@ export function startHttpServer() {
res.setHeader("Cache-Control", "no-store"); res.setHeader("Cache-Control", "no-store");
res.type(attachment.type).send(Buffer.from(data, "base64")); res.type(attachment.type).send(Buffer.from(data, "base64"));
})); }));
app.get("/agent/message-assets/:messageKey/:assetFile", route(async (req, res) => {
const asset = await messageMetadataStore.readAsset(routeParam(req.params.messageKey), routeParam(req.params.assetFile));
if (!asset) return void res.status(404).json({ ok: false, error: "message asset not found" });
res.setHeader("Cache-Control", "private, max-age=31536000, immutable");
res.type(asset.contentType).send(asset.data);
}));
app.post("/agent/local-file/reveal", route(async (req, res) => { app.post("/agent/local-file/reveal", route(async (req, res) => {
const filePath = String(req.body?.path || ""); const filePath = String(req.body?.path || "");
if (!path.isAbsolute(filePath)) return res.status(400).json({ ok: false, error: "文件路径必须是绝对路径" }); if (!path.isAbsolute(filePath)) return res.status(400).json({ ok: false, error: "文件路径必须是绝对路径" });
@@ -305,6 +312,7 @@ export function startHttpServer() {
const skill = req.body?.skill === undefined ? undefined : await resolveCodexSkill(emit, workspace.workspacePath, skillSelector(req.body.skill), true); const skill = req.body?.skill === undefined ? undefined : await resolveCodexSkill(emit, workspace.workspacePath, skillSelector(req.body.skill), true);
const messageId = String(req.body?.messageId || Date.now()); const messageId = String(req.body?.messageId || Date.now());
const messageText = String(req.body?.messageText || prompt || `发送了 ${attachments.length} 张图片`); const messageText = String(req.body?.messageText || prompt || `发送了 ${attachments.length} 张图片`);
const messageMetadata = await messageMetadataStore.recordPending(messageId, req.body?.messageMetadata);
let threadId = activeThreadId; let threadId = activeThreadId;
logger.info("Codex turn accepted", { threadId: req.body?.threadId, model: model || "default", reasoningEffort: effort || "default", promptLength: prompt.length, attachmentCount: attachments.length }); logger.info("Codex turn accepted", { threadId: req.body?.threadId, model: model || "default", reasoningEffort: effort || "default", promptLength: prompt.length, attachmentCount: attachments.length });
session.bindClient(clientId); session.bindClient(clientId);
@@ -315,7 +323,7 @@ export function startHttpServer() {
const attachmentRefs = session.setTurnAttachments(clientId, attachments); const attachmentRefs = session.setTurnAttachments(clientId, attachments);
session.emitThread("chat_message", threadId, { session.emitThread("chat_message", threadId, {
sourceClientId: clientId, sourceClientId: clientId,
message: { id: `${threadId}:pending:synthetic:user`, itemId: "synthetic:user", clientMessageId: messageId, threadId, turnId: "", role: "user", text: messageText }, message: { id: `${threadId}:pending:synthetic:user`, itemId: "synthetic:user", clientMessageId: messageId, threadId, turnId: "", role: "user", text: messageText, ...messageMetadata },
}); });
let chatTurnId = ""; let chatTurnId = "";
/** 将包装层日志和兜底错误固定广播到当前 turn。 */ /** 将包装层日志和兜底错误固定广播到当前 turn。 */
@@ -344,6 +352,7 @@ export function startHttpServer() {
onStart: () => session.bindClient(clientId), onStart: () => session.bindClient(clientId),
onThread: (actualThreadId) => { onThread: (actualThreadId) => {
const threadChanged = actualThreadId !== threadId; const threadChanged = actualThreadId !== threadId;
void messageMetadataStore.bindThread(messageId, actualThreadId).catch((error) => logger.warn("Failed to bind message metadata to thread", { clientMessageId: messageId, threadId: actualThreadId, error }));
if (actualThreadId !== threadId) { if (actualThreadId !== threadId) {
threadId = actualThreadId; threadId = actualThreadId;
setActiveThread(threadId, { emptyThread: true, sourceClientId: clientId }); setActiveThread(threadId, { emptyThread: true, sourceClientId: clientId });
@@ -353,18 +362,19 @@ export function startHttpServer() {
if (threadChanged) { if (threadChanged) {
session.emitThread("chat_message", threadId, { session.emitThread("chat_message", threadId, {
sourceClientId: clientId, sourceClientId: clientId,
message: { id: `${threadId}:pending:synthetic:user`, itemId: "synthetic:user", clientMessageId: messageId, threadId, turnId: "", role: "user", text: messageText }, message: { id: `${threadId}:pending:synthetic:user`, itemId: "synthetic:user", clientMessageId: messageId, threadId, turnId: "", role: "user", text: messageText, ...messageMetadata },
}); });
} }
}, },
onTurn: (actualTurnId) => { onTurn: (actualTurnId) => {
turnId = actualTurnId; turnId = actualTurnId;
void messageMetadataStore.bindTurn(messageId, threadId, turnId).catch((error) => logger.warn("Failed to bind message metadata to turn", { clientMessageId: messageId, threadId, turnId, error }));
if (chatTurnId !== turnId) { if (chatTurnId !== turnId) {
chatTurnId = turnId; chatTurnId = turnId;
session.emitThread("chat_message", threadId, { session.emitThread("chat_message", threadId, {
turnId, turnId,
sourceClientId: clientId, sourceClientId: clientId,
message: { id: `${threadId}:${turnId}:synthetic:user`, itemId: "synthetic:user", clientMessageId: messageId, threadId, turnId, role: "user", text: messageText }, message: { id: `${threadId}:${turnId}:synthetic:user`, itemId: "synthetic:user", clientMessageId: messageId, threadId, turnId, role: "user", text: messageText, ...messageMetadata },
}); });
} }
logger.info("Codex turn started", { threadId, turnId, model: model || "default", reasoningEffort: effort || "default" }); logger.info("Codex turn started", { threadId, turnId, model: model || "default", reasoningEffort: effort || "default" });
@@ -372,6 +382,7 @@ export function startHttpServer() {
}, },
onFinish: () => { onFinish: () => {
logger.info("Codex turn finished", { threadId, turnId }); logger.info("Codex turn finished", { threadId, turnId });
if (!turnId) void messageMetadataStore.remove(messageId, threadId).catch((error) => logger.warn("Failed to remove unbound message metadata", { clientMessageId: messageId, error }));
session.clearTurnAttachments(clientId); session.clearTurnAttachments(clientId);
if (clientId) session.releaseClient(clientId); if (clientId) session.releaseClient(clientId);
session.setCodexState({ busy: false, threadId, turnId }); session.setCodexState({ busy: false, threadId, turnId });
@@ -380,6 +391,7 @@ export function startHttpServer() {
}); });
res.json({ ok: true, threadId }); res.json({ ok: true, threadId });
} catch (error) { } catch (error) {
await messageMetadataStore.remove(messageId, threadId).catch((metadataError) => logger.warn("Failed to remove rejected message metadata", { clientMessageId: messageId, error: metadataError }));
session.releaseClient(clientId); session.releaseClient(clientId);
session.setCodexState({ busy: false, threadId, turnId: "" }); session.setCodexState({ busy: false, threadId, turnId: "" });
session.finishConversationRun(threadId); session.finishConversationRun(threadId);
@@ -71,7 +71,7 @@ const MAX_ATTACHMENT_PAYLOAD_BYTES = 28 * 1024 * 1024;
const MESSAGE_PREVIEW_LONG_EDGE = 192; const MESSAGE_PREVIEW_LONG_EDGE = 192;
const MESSAGE_PREVIEW_MAX_LENGTH = 500_000; const MESSAGE_PREVIEW_MAX_LENGTH = 500_000;
const DEFAULT_AGENT_URL = "http://127.0.0.1:17371"; const DEFAULT_AGENT_URL = "http://127.0.0.1:17371";
const AGENT_PROTOCOL_VERSION = 5; const AGENT_PROTOCOL_VERSION = 6;
const HISTORY_RETRY_DELAYS_MS = [0, 150, 350, 700, 1200]; const HISTORY_RETRY_DELAYS_MS = [0, 150, 350, 700, 1200];
const AGENT_REASONING_EFFORTS = new Set<AgentReasoningEffort>(["minimal", "low", "medium", "high", "xhigh", "max", "ultra"]); const AGENT_REASONING_EFFORTS = new Set<AgentReasoningEffort>(["minimal", "low", "medium", "high", "xhigh", "max", "ultra"]);
const rt = (key: string, options?: Record<string, unknown>) => i18n.t(`agent.runtime.${key}`, options); const rt = (key: string, options?: Record<string, unknown>) => i18n.t(`agent.runtime.${key}`, options);
-132
View File
@@ -1,132 +0,0 @@
import localforage from "localforage";
import { upscaleDataUrl } from "@/lib/canvas/canvas-image-data";
import type { AgentAttachment, AgentChatItem } from "@/stores/use-agent-store";
export type StoredAgentUserMessage = Pick<AgentChatItem, "id" | "text" | "attachments"> & { role: "user"; historyText: string; threadId?: string; turnId?: string };
const store = localforage.createInstance({ name: "infinite-canvas", storeName: "agent_chat_messages" });
const mutations = new Map<string, Promise<void>>();
const indexKey = (threadId: string) => `thread:${threadId}`;
const messageKey = (threadId: string, messageId: string) => `message:${threadId}:${messageId}`;
const pendingKey = (messageId: string) => `pending:${messageId}`;
const threadMutationKey = (threadId: string) => `thread:${threadId}`;
const pendingMutationKey = (messageId: string) => `pending:${messageId}`;
export async function saveAgentUserMessage(threadId: string, message: StoredAgentUserMessage) {
if (!message.attachments?.length) return;
if (!threadId) return savePendingAgentUserMessage(message);
await saveThreadAgentUserMessage(threadId, message);
}
/** Persist attachments before a turn is accepted. The record is moved to a thread after the server assigns one. */
export async function savePendingAgentUserMessage(message: StoredAgentUserMessage) {
if (!message.id || !message.attachments?.length) return;
await mutateScopes([pendingMutationKey(message.id)], async () => {
const attachments = await Promise.all(message.attachments!.map(createThumbnail));
await store.setItem(pendingKey(message.id), { ...message, threadId: undefined, turnId: undefined, attachments });
});
}
export async function deletePendingAgentUserMessage(messageId: string) {
if (!messageId) return;
await mutateScopes([pendingMutationKey(messageId)], () => store.removeItem(pendingKey(messageId)));
}
export async function readAgentUserMessages(threadId: string) {
await mutations.get(threadMutationKey(threadId))?.catch(() => undefined);
const ids = (await store.getItem<string[]>(indexKey(threadId))) || [];
return (await Promise.all(ids.map((id) => store.getItem<StoredAgentUserMessage>(messageKey(threadId, id))))).filter((item): item is StoredAgentUserMessage => Boolean(item));
}
/** Bind a pending message to the server thread, preserving an already-known turn id. */
export async function bindPendingAgentUserMessage(threadId: string, messageId: string, turnId = "") {
if (!threadId || !messageId) return;
await mutateScopes([pendingMutationKey(messageId), threadMutationKey(threadId)], async () => {
const pending = await store.getItem<StoredAgentUserMessage>(pendingKey(messageId));
const key = messageKey(threadId, messageId);
const existing = await store.getItem<StoredAgentUserMessage>(key);
if (!pending && !existing) return;
const message = mergeStoredMessage(existing, pending, threadId, turnId);
await putThreadMessage(threadId, key, message);
if (pending) await store.removeItem(pendingKey(messageId));
});
}
export async function bindAgentUserMessageTurn(threadId: string, messageId: string, turnId: string) {
await bindPendingAgentUserMessage(threadId, messageId, turnId);
}
export async function moveAgentUserMessage(fromThreadId: string, toThreadId: string, messageId: string) {
if (!toThreadId || !messageId || fromThreadId === toThreadId) return bindPendingAgentUserMessage(toThreadId, messageId);
const scopes = [pendingMutationKey(messageId), threadMutationKey(toThreadId), ...(fromThreadId ? [threadMutationKey(fromThreadId)] : [])];
await mutateScopes(scopes, async () => {
const pending = await store.getItem<StoredAgentUserMessage>(pendingKey(messageId));
const fromKey = fromThreadId ? messageKey(fromThreadId, messageId) : "";
const from = fromKey ? await store.getItem<StoredAgentUserMessage>(fromKey) : null;
const toKey = messageKey(toThreadId, messageId);
const existing = await store.getItem<StoredAgentUserMessage>(toKey);
const source = pending || from;
if (!source && !existing) return;
await putThreadMessage(toThreadId, toKey, mergeStoredMessage(existing, source, toThreadId));
if (pending) await store.removeItem(pendingKey(messageId));
if (from && fromThreadId) await removeThreadMessage(fromThreadId, fromKey, messageId);
});
}
export async function deleteAgentThreadMessages(threadIds: string[]) {
await mutateScopes(threadIds.map(threadMutationKey), async () => {
await Promise.all(threadIds.map(async (threadId) => {
const ids = (await store.getItem<string[]>(indexKey(threadId))) || [];
await Promise.all(ids.map((id) => store.removeItem(messageKey(threadId, id))));
await store.removeItem(indexKey(threadId));
}));
});
}
async function saveThreadAgentUserMessage(threadId: string, message: StoredAgentUserMessage) {
await mutateScopes([threadMutationKey(threadId)], async () => {
const attachments = await Promise.all(message.attachments!.map(createThumbnail));
await putThreadMessage(threadId, messageKey(threadId, message.id), { ...message, threadId, attachments });
});
}
async function putThreadMessage(threadId: string, key: string, message: StoredAgentUserMessage) {
await store.setItem(key, { ...message, threadId });
const ids = (await store.getItem<string[]>(indexKey(threadId))) || [];
if (!ids.includes(message.id)) await store.setItem(indexKey(threadId), [...ids, message.id]);
}
function mergeStoredMessage(existing: StoredAgentUserMessage | null, source: StoredAgentUserMessage | null | undefined, threadId: string, turnId = "") {
const message = { ...(source || {}), ...(existing || {}) } as StoredAgentUserMessage;
if (!message.attachments?.length && source?.attachments?.length) message.attachments = source.attachments;
if (!message.text && source?.text) message.text = source.text;
if (!message.historyText && source?.historyText) message.historyText = source.historyText;
return { ...message, threadId, ...(turnId ? { turnId } : message.turnId ? { turnId: message.turnId } : {}) };
}
async function removeThreadMessage(threadId: string, key: string, messageId: string) {
await store.removeItem(key);
const ids = (await store.getItem<string[]>(indexKey(threadId))) || [];
const remaining = ids.filter((id) => id !== messageId);
if (remaining.length) await store.setItem(indexKey(threadId), remaining);
else await store.removeItem(indexKey(threadId));
}
async function mutateScopes(scopes: string[], mutation: () => Promise<void>) {
const ids = [...new Set(scopes.filter(Boolean))].sort();
const operation = Promise.all(ids.map((id) => mutations.get(id)?.catch(() => undefined))).then(mutation);
ids.forEach((id) => mutations.set(id, operation));
try {
await operation;
} finally {
ids.forEach((id) => {
if (mutations.get(id) === operation) mutations.delete(id);
});
}
}
async function createThumbnail(attachment: AgentAttachment): Promise<AgentAttachment> {
const dataUrl = Math.max(attachment.width, attachment.height) > 512 ? await upscaleDataUrl(attachment.dataUrl, { targetLongEdge: 512, algorithm: "high" }) : attachment.dataUrl;
return { ...attachment, size: dataUrl.length, url: dataUrl, dataUrl };
}