Files
MiragenFlow/web/src/components/agent/local-agent-panel.tsx
T

1468 lines
82 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 { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { useNavigate, useSearchParams } from "react-router-dom";
import { App, Button, Tooltip } from "antd";
import dayjs from "dayjs";
import { Bot, History, MessageSquare, PanelRightClose, PlugZap, Plus, Terminal } from "lucide-react";
import { canvasThemes } from "@/lib/canvas-theme";
import { imageMetadata } from "@/lib/canvas/canvas-node-factory";
import { fitNodeSize } from "@/lib/canvas/canvas-node-size";
import { readImageMeta } from "@/lib/image-utils";
import { randomId } from "@/lib/utils";
import { uploadImage } from "@/services/image-storage";
import { bindPendingAgentUserMessage, deleteAgentThreadMessages, deletePendingAgentUserMessage, moveAgentUserMessage, readAgentUserMessages, savePendingAgentUserMessage } from "@/services/agent-chat-storage";
import { useThemeStore } from "@/stores/use-theme-store";
import { useShallow } from "zustand/react/shallow";
import { useAgentStore, type AgentCanvasContext, type AgentChatItem, type AgentModel, type AgentPendingApproval, type AgentPendingToolCall, type AgentPermissionMode, type AgentReasoningEffort, type AgentThreadSummary } from "@/stores/use-agent-store";
import { type CanvasAgentOp, type CanvasAgentSnapshot } from "@/lib/canvas/canvas-agent-ops";
import { isSiteTool, runSiteTool } from "@/lib/agent/agent-site-tools";
import { acknowledgeCodexHistory, activateAgentClient, discoverAgentConfig, fetchAgentJson, postCodexApproval, postState, postToolResult } from "./agent-api";
import { AgentChatTimeline, AgentTaskProgress, AgentUsageBar } from "./agent-chat";
import { AgentChatComposer } from "./agent-chat-composer";
import { AgentConnectView } from "./agent-connect-view";
import {
activityDeltaFallback,
activityDetail,
activityKind,
activityPlaceholder,
agentAttachmentToChatAttachment,
agentErrorView,
attachmentPayloadBytes,
compactText,
eventUsage,
formatAgentActivity,
formatAgentEvent,
formatAgentEventLog,
formatAgentPlan,
formatBytes,
bindPendingTurnMessages,
isCanvasWriteTool,
isConnectionErrorMessage,
isCurrentThreadEvent,
isReasoningSummary,
mergeAgentMessages,
mergeHistoryAttachments,
mergeStreamText,
normalizeHistoryMessages,
normalizeText,
parseEventData,
promptWithAttachments,
reasoningActivityText,
registerLiveAgentTurn,
scopeChatItem,
stringText,
toolName,
turnPlanStatus,
upsertAgentMessage,
type AgentEventItem,
type AgentEventPayload,
} from "./agent-event-formatters";
import { AgentHistoryView } from "./agent-history-view";
import { AgentLogView } from "./agent-log-view";
import { AgentPanelTabs } from "./agent-panel-tabs";
const MAX_ATTACHMENTS = 6;
const MAX_ATTACHMENT_PAYLOAD_BYTES = 28 * 1024 * 1024;
const DEFAULT_AGENT_URL = "http://127.0.0.1:17371";
const AGENT_PROTOCOL_VERSION = 3;
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_LABELS: Record<AgentReasoningEffort, string> = { minimal: "最低", low: "轻度", medium: "中", high: "高", xhigh: "极高", max: "最高", ultra: "Ultra" };
type AgentWorkspace = { workspacePath: string; activeThreadId?: string };
type AgentThreadsResponse = { ok?: boolean; workspace?: AgentWorkspace; data?: AgentThreadSummary[] };
type AgentThreadResponse = { ok?: boolean; workspace?: AgentWorkspace; thread?: AgentThreadSummary; messages?: AgentChatItem[]; settledTurnIds?: string[]; historyReady?: boolean };
type AgentWorkspaceResponse = { ok?: boolean; workspace?: AgentWorkspace };
type AgentTurnResponse = { ok?: boolean; threadId?: string };
type AgentModelsResponse = { ok?: boolean; data?: AgentModel[] };
type AgentCodexState = { busy?: boolean; threadId?: string; turnId?: string };
type AgentHelloEvent = { ok?: boolean; protocolVersion?: number; clientId?: string; workspace?: { activeThreadId?: string }; codex?: AgentCodexState; pendingApprovals?: AgentPendingApproval[] };
type AgentWorkspaceEvent = { activeThreadId?: string; threadId?: string; sourceClientId?: string; emptyThread?: boolean; draftThread?: boolean };
type AgentChatEvent = { threadId?: string; turnId?: string; sourceClientId?: string; replayed?: boolean; message?: AgentChatItem };
type AgentBootstrapEvent = { type?: "codex.preparing" | "codex.prepare_failed" | "mcp.startup"; threadId?: string; name?: string; status?: "starting" | "ready" | "failed" | "cancelled"; error?: string | null; failureReason?: string | null };
type AgentClientGlobal = typeof globalThis & { __infiniteCanvasAgentClientIdPromise?: Promise<string> };
function authoritativeHistoryTurnKeys(threadId: string, settledTurnIds: string[]) {
return new Set(settledTurnIds.map((turnId) => `${threadId}\0${turnId}`));
}
export function LocalAgentPanel({ embedded, headless, autoConnect }: { embedded?: boolean; headless?: boolean; autoConnect?: boolean }) {
const theme = canvasThemes[useThemeStore((state) => state.theme)];
const { message, modal } = App.useApp();
const [searchParams] = useSearchParams();
const navigate = useNavigate();
// 逐字段 selector + useShallow:只有这些字段变化时才重渲染。
// 注意:canvasContext 不在此订阅内 —— 它在拖拽/resize 时会被 project 每帧写入,
// 但面板只在 ref 同步与防抖 postState 中用到它、渲染层从不读它。若把它放进订阅,
// 面板会随画布每帧重渲染(性能问题,也是 #185 崩溃的放大器)。改为下方 subscribe 命令式监听。
const { width, url, token, connected, enabled, prompt, attachments, sending, waiting, tokenUsage, eventLogs, threads, activeThreadId, workspacePath, loadingThreads, activeTab, confirmTools, permissionMode, models, model, reasoningEffort, activity, connectError, pendingTool, pendingApprovals } = useAgentStore(
useShallow((state) => ({
width: state.width,
url: state.url,
token: state.token,
connected: state.connected,
enabled: state.enabled,
prompt: state.prompt,
attachments: state.attachments,
sending: state.sending,
waiting: state.waiting,
tokenUsage: state.tokenUsage,
eventLogs: state.eventLogs,
threads: state.threads,
activeThreadId: state.activeThreadId,
workspacePath: state.workspacePath,
loadingThreads: state.loadingThreads,
activeTab: state.activeTab,
confirmTools: state.confirmTools,
permissionMode: state.permissionMode,
models: state.models,
model: state.model,
reasoningEffort: state.reasoningEffort,
activity: state.activity,
connectError: state.connectError,
pendingTool: state.pendingTool,
pendingApprovals: state.pendingApprovals,
})),
);
const setAgentState = useAgentStore((state) => state.setAgentState);
const agentInitializing = useAgentStore((state) => state.bootstrapStatus?.status === "running");
const closePanel = useAgentStore((state) => state.closePanel);
const pushMessage = useAgentStore((state) => state.addMessage);
const pushEventLog = useAgentStore((state) => state.addEventLog);
const clearEventLogs = useAgentStore((state) => state.clearEventLogs);
const messageCount = useAgentStore((state) => state.messages.length);
const canvasContextRef = useRef<AgentCanvasContext | null>(useAgentStore.getState().canvasContext);
const confirmToolsRef = useRef(confirmTools);
const pendingToolRef = useRef<AgentPendingToolCall | null>(null);
const autoConnectRef = useRef(false);
const connectedRef = useRef(false);
const errorLoggedRef = useRef(false);
const attachmentUrlsRef = useRef(new Set<string>());
const clientIdRef = useRef("");
const [clientReady, setClientReady] = useState(false);
const loadThreadsSequenceRef = useRef(0);
const threadMessagesRef = useRef(new Map<string, AgentChatItem[]>());
const authoritativeHistoryTurnsRef = useRef(new Set<string>());
const liveTurnKeysRef = useRef(new Set<string>());
const threadOperationRef = useRef(0);
const threadOperationSequenceRef = useRef(0);
const endpoint = useMemo(() => url.trim().replace(/\/$/, ""), [url]);
const urlAgentAutoConnect = searchParams.has("agentUrl") && searchParams.has("agentToken");
useEffect(() => {
let disposed = false;
void acquireAgentClientId().then((clientId) => {
if (!disposed) {
clientIdRef.current = clientId;
setClientReady(true);
}
});
return () => { disposed = true; };
}, []);
const loadThreadSnapshot = useCallback(async (threadId: string, sequence: number, response?: AgentThreadResponse, expectedTurnId = "") => {
const storedMessagesPromise = readAgentUserMessages(threadId).catch(() => []);
let thread = response;
let lastError: unknown;
for (const delayMs of HISTORY_RETRY_DELAYS_MS) {
if (delayMs) await delay(delayMs);
if (sequence !== loadThreadsSequenceRef.current || useAgentStore.getState().activeThreadId !== threadId) return false;
try {
thread ||= await fetchAgentJson<AgentThreadResponse>(endpoint, token, `/agent/codex/threads/${encodeURIComponent(threadId)}`);
lastError = undefined;
} catch (error) {
lastError = error;
thread = undefined;
continue;
}
const history = normalizeHistoryMessages(thread.messages || []);
const storedMessages = await storedMessagesPromise;
const latest = useAgentStore.getState();
if (sequence !== loadThreadsSequenceRef.current || latest.activeThreadId !== threadId) return false;
const historyTurns = authoritativeHistoryTurnKeys(threadId, thread.settledTurnIds || []);
const hasExpectedTurn = !expectedTurnId || historyTurns.has(`${threadId}\0${expectedTurnId}`);
historyTurns.forEach((key) => liveTurnKeysRef.current.delete(key));
if (latest.activeTurnId) liveTurnKeysRef.current.add(`${threadId}\0${latest.activeTurnId}`);
authoritativeHistoryTurnsRef.current = historyTurns;
const attachmentSources = [...latest.messages, ...storedMessages.map((item) => ({ ...item, threadId }))];
const snapshot = mergeHistoryAttachments(history, attachmentSources);
const messages = mergeAgentMessages(snapshot, latest.messages, threadId, liveTurnKeysRef.current);
threadMessagesRef.current.set(threadId, messages);
setAgentState({ messages, connectError: "" });
const coveredTurnIds = [...historyTurns].map((key) => key.slice(threadId.length + 1));
if (coveredTurnIds.length) void acknowledgeCodexHistory(endpoint, token, threadId, coveredTurnIds).catch(() => undefined);
if (hasExpectedTurn && (thread.historyReady !== false || Boolean(expectedTurnId))) return true;
thread = undefined;
}
if (lastError) throw lastError;
return false;
}, [endpoint, setAgentState, token]);
const applyWorkspaceChange = useCallback((data: AgentWorkspaceEvent) => {
const nextThreadId = data.activeThreadId ?? data.threadId ?? "";
const current = useAgentStore.getState();
const threadChanged = current.activeThreadId !== nextThreadId;
const emptyThread = Boolean(data.emptyThread || data.draftThread);
const pendingMessage = [...current.messages].reverse().find((item) => item.role === "user" && !item.turnId);
const keepPendingMessage = Boolean(
data.emptyThread
&& pendingMessage
&& (current.sending || current.waiting)
&& (!data.sourceClientId || data.sourceClientId === clientIdRef.current),
);
if (threadChanged && current.activeThreadId) {
const messages = keepPendingMessage ? current.messages.filter((item) => item.id !== pendingMessage!.id) : current.messages;
threadMessagesRef.current.set(current.activeThreadId, messages);
}
if (emptyThread && nextThreadId) threadMessagesRef.current.delete(nextThreadId);
if (threadChanged || emptyThread) {
loadThreadsSequenceRef.current += 1;
authoritativeHistoryTurnsRef.current.clear();
liveTurnKeysRef.current.clear();
}
const messages = keepPendingMessage
? [scopeChatItem(pendingMessage!, nextThreadId, "")]
: emptyThread ? []
: threadChanged ? threadMessagesRef.current.get(nextThreadId) || []
: current.messages;
pendingToolRef.current = null;
setAgentState({
activeThreadId: nextThreadId,
activeTurnId: threadChanged || emptyThread ? "" : current.activeTurnId,
messages,
tokenUsage: threadChanged || emptyThread ? null : current.tokenUsage,
pendingTool: null,
pendingApprovals: threadChanged || emptyThread ? [] : current.pendingApprovals,
});
return loadThreadsSequenceRef.current;
}, [setAgentState]);
const loadThreads = useCallback(async (skipHistory = false, expectedTurnId = "") => {
if (!connectedRef.current && !useAgentStore.getState().connected) return;
let sequence = ++loadThreadsSequenceRef.current;
setAgentState({ loadingThreads: true });
try {
const data = await fetchAgentJson<AgentThreadsResponse>(endpoint, token, `/agent/codex/threads`);
if (sequence !== loadThreadsSequenceRef.current) return;
const current = useAgentStore.getState();
const currentThreadId = data.workspace?.activeThreadId ?? current.activeThreadId;
if (currentThreadId !== current.activeThreadId) sequence = applyWorkspaceChange({ activeThreadId: currentThreadId });
if (sequence !== loadThreadsSequenceRef.current || useAgentStore.getState().activeThreadId !== currentThreadId) return;
setAgentState({ threads: data.data || [], workspacePath: data.workspace?.workspacePath || "" });
if (currentThreadId && !skipHistory) {
await loadThreadSnapshot(currentThreadId, sequence, undefined, expectedTurnId);
} else {
authoritativeHistoryTurnsRef.current.clear();
liveTurnKeysRef.current.clear();
}
} catch (error) {
addEventLog("读取历史失败", error);
} finally {
if (sequence === loadThreadsSequenceRef.current && !threadOperationRef.current) setAgentState({ loadingThreads: false });
}
}, [applyWorkspaceChange, endpoint, loadThreadSnapshot, setAgentState, token]);
// canvasContext 命令式订阅:保持 ref 最新,并在快照变化时防抖上报,全程不触发面板重渲染。
useEffect(() => {
let timer: ReturnType<typeof setTimeout> | null = null;
const unsubscribe = useAgentStore.subscribe((state) => {
if (state.canvasContext === canvasContextRef.current) return;
canvasContextRef.current = state.canvasContext;
if (!useAgentStore.getState().connected) return;
if (timer) clearTimeout(timer);
timer = setTimeout(() => void postState(endpoint, token, clientIdRef.current, canvasContextRef.current?.snapshot || null), 300);
});
return () => {
unsubscribe();
if (timer) clearTimeout(timer);
};
}, [endpoint, token]);
useEffect(() => {
confirmToolsRef.current = confirmTools;
}, [confirmTools]);
useEffect(() => {
pendingToolRef.current = pendingTool;
}, [pendingTool]);
useEffect(() => () => attachmentUrlsRef.current.forEach((url) => URL.revokeObjectURL(url)), []);
useEffect(() => {
if (!clientReady || !enabled || !token.trim()) return;
localStorage.setItem("canvas-agent-url", endpoint);
localStorage.setItem("canvas-agent-token", token);
const clientId = clientIdRef.current;
let disposed = false;
let protocolRejected = false;
let eventQueue = Promise.resolve();
const isCurrentConnection = () => !disposed && clientIdRef.current === clientId;
const enqueueEvent = (task: () => void | Promise<void>) => {
eventQueue = eventQueue.then(async () => {
if (isCurrentConnection()) await task();
}).catch((error) => {
if (isCurrentConnection()) addEventLog("同步会话失败", error);
});
};
const source = new EventSource(`${endpoint}/events?token=${encodeURIComponent(token)}&clientId=${encodeURIComponent(clientId)}`);
source.addEventListener("hello", (event) => {
if (!isCurrentConnection()) return;
const hello = parseEventData<AgentHelloEvent>(event);
if (hello?.protocolVersion !== AGENT_PROTOCOL_VERSION) {
const text = "本地 Agent 版本过旧,请重启 Canvas Agent 后重新连接";
protocolRejected = true;
source.close();
connectedRef.current = false;
setAgentState({ enabled: false, connected: false, waiting: false, sending: false, activity: "需要重启 Agent", connectError: text, silentConnect: false, pendingTool: null, pendingApprovals: [] });
addEventLog("Agent 版本不匹配", text, hello);
if (!headless) message.error(text);
return;
}
const codex = hello?.codex;
const busy = Boolean(codex?.busy);
const nextThreadId = hello?.workspace?.activeThreadId ?? useAgentStore.getState().activeThreadId;
applyWorkspaceChange({ activeThreadId: nextThreadId });
const current = useAgentStore.getState();
const nextTurnId = codex?.threadId === nextThreadId ? codex.turnId ?? "" : "";
if (nextTurnId) liveTurnKeysRef.current.add(`${nextThreadId}\0${nextTurnId}`);
const activeTurnId = busy ? nextTurnId : "";
const pendingApprovals = busy ? (hello?.pendingApprovals || []).filter((item) => !item.threadId || item.threadId === nextThreadId) : [];
const messages = activeTurnId
? bindPendingTurnMessages(current.messages.filter((item) => !isConnectionErrorMessage(item)), nextThreadId, activeTurnId)
: current.messages.filter((item) => !isConnectionErrorMessage(item));
errorLoggedRef.current = false;
connectedRef.current = true;
setAgentState({
connected: true,
activity: pendingApprovals.length ? "等待权限确认" : busy ? "Codex 正在运行" : "已连接",
waiting: busy,
sending: false,
connectError: "",
silentConnect: false,
activeThreadId: nextThreadId,
activeTurnId,
messages,
pendingApprovals,
});
if (!headless) message.success("本地 Agent 已连接");
void postState(endpoint, token, clientId, canvasContextRef.current?.snapshot || null);
if (document.visibilityState === "visible" && document.hasFocus()) void activateAgentClient(endpoint, token, clientId);
if (!busy && !nextThreadId) {
setAgentState({ bootstrapStatus: { key: "codex:preparing", text: "正在初始化 Codex 对话", detail: "正在创建会话并启动画布工具服务", status: "running" }, mcpStartupStatuses: {} });
void fetchAgentJson(endpoint, token, "/agent/codex/threads/reset", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ clientId, permissionMode }) }).catch((error) => {
setAgentState({ bootstrapStatus: { key: "codex:prepare_failed", text: "Codex 对话初始化失败", detail: error instanceof Error ? error.message : "无法创建 Codex 会话", status: "error" } });
addEventLog("Codex 对话初始化失败", error);
});
}
});
source.addEventListener("codex_state", (event) => {
const data = parseEventData<AgentCodexState>(event);
if (!data) return;
enqueueEvent(async () => {
const busy = Boolean(data.busy);
const current = useAgentStore.getState();
const appliesToCurrentThread = !data.threadId || data.threadId === current.activeThreadId;
if (!appliesToCurrentThread) return;
const turnId = data.turnId || current.activeTurnId;
if (turnId) liveTurnKeysRef.current.add(`${current.activeThreadId}\0${turnId}`);
const activeTurnId = busy ? turnId : "";
const messages = activeTurnId ? bindPendingTurnMessages(current.messages, current.activeThreadId, activeTurnId) : current.messages;
setAgentState({
activity: busy ? "Codex 正在运行" : current.activity === "处理失败" ? "处理失败" : "完成",
waiting: busy,
sending: false,
activeTurnId,
messages,
});
if (!busy && current.waiting) void loadThreads(false, turnId);
});
});
source.addEventListener("tool_call", (event) => {
if (!isCurrentConnection()) return;
const data = parseEventData<AgentPendingToolCall>(event);
if (data) void handleToolCall(endpoint, token, data);
});
source.addEventListener("codex_approval", (event) => {
if (!isCurrentConnection()) return;
const data = parseEventData<AgentPendingApproval>(event);
if (!data || !isCurrentThreadEvent(data)) return;
setAgentState({ pendingApprovals: [...useAgentStore.getState().pendingApprovals.filter((item) => item.requestId !== data.requestId), data], activity: "等待权限确认" });
addEventLog("等待权限确认", data.reason || data.method, data);
});
source.addEventListener("codex_approval_resolved", (event) => {
if (!isCurrentConnection()) return;
const data = parseEventData<{ requestId?: string; decision?: "accept" | "acceptForSession" | "decline" | "cancel" }>(event);
if (!data?.requestId) return;
const current = useAgentStore.getState();
const approval = current.pendingApprovals.find((item) => item.requestId === data.requestId);
const pendingApprovals = current.pendingApprovals.filter((item) => item.requestId !== data.requestId);
setAgentState({ pendingApprovals, activity: approvalActivity(pendingApprovals, current.waiting, current.activity) });
const decision = data.decision || approval?.deciding;
if (approval && decision) addEventLog(decision === "accept" || decision === "acceptForSession" ? "已批准权限" : "已取消权限", approval.reason || approval.method, approval);
});
source.addEventListener("agent_event", (event) => {
const data = parseEventData<AgentEventPayload>(event);
if (data) enqueueEvent(() => {
if (!isCurrentThreadEvent(data)) return;
const shouldProcess = registerLiveAgentTurn(data, authoritativeHistoryTurnsRef.current, liveTurnKeysRef.current);
if (data.type !== "usage.updated" && !shouldProcess) return;
return handleAgentEvent(data);
});
});
source.addEventListener("agent_bootstrap", (event) => {
const data = parseEventData<AgentBootstrapEvent>(event);
if (!data?.type) return;
if (data.type === "codex.preparing") {
setAgentState({ bootstrapStatus: { key: "codex:preparing", text: "正在初始化 Codex 对话", detail: "正在创建会话并启动画布工具服务", status: "running" }, mcpStartupStatuses: {} });
addEventLog("正在初始化 Codex 对话", "正在创建会话并启动画布工具服务", data);
return;
}
if (data.type === "codex.prepare_failed") {
setAgentState({ bootstrapStatus: { key: "codex:prepare_failed", text: "Codex 对话初始化失败", detail: data.error || "无法创建 Codex 会话", status: "error" } });
addEventLog("Codex 对话初始化失败", data.error, data);
return;
}
if (!data.name || !data.status) return;
const label = data.name;
const status = data.status === "starting"
? { text: `正在启动 MCP:${label}`, detail: "正在建立工具连接并读取可用工具列表", status: "running" as const }
: data.status === "ready"
? { text: `MCP 已就绪:${label}`, detail: "工具列表加载完成,可以开始对话", status: "ready" as const }
: data.status === "failed"
? { text: `MCP 启动失败:${label}`, detail: data.error || "工具服务未能完成初始化", status: "error" as const }
: { text: `MCP 启动已取消:${label}`, detail: "工具服务初始化已取消", status: "error" as const };
const mcpStartupStatuses = { ...useAgentStore.getState().mcpStartupStatuses, [label]: { key: `mcp:${label}:${data.status}`, ...status } };
const services = Object.values(mcpStartupStatuses);
const failed = services.some((item) => item.status === "error");
const ready = services.length > 0 && services.every((item) => item.status === "ready");
setAgentState({
mcpStartupStatuses,
bootstrapStatus: failed
? { key: "mcp:failed", text: "部分 MCP 服务初始化失败", detail: "可以查看下方服务状态和诊断日志", status: "error" }
: ready
? { key: "mcp:ready", text: `${services.length} 个 MCP 服务已就绪`, detail: "工具列表加载完成,可以开始对话", status: "ready" }
: { key: "mcp:starting", text: "正在启动 MCP 服务", detail: `正在初始化 ${services.length} 个工具服务`, status: "running" },
});
addEventLog(status.text, status.detail, data);
});
source.addEventListener("workspace_changed", (event) => {
const data = parseEventData<AgentWorkspaceEvent>(event);
if (!data) return;
enqueueEvent(async () => {
const nextThreadId = data.activeThreadId ?? data.threadId ?? "";
const current = useAgentStore.getState();
const pendingMessage = [...current.messages].reverse().find((item) => item.role === "user" && !item.turnId);
const keepPendingMessage = Boolean(data.emptyThread && pendingMessage && (current.sending || current.waiting) && (!data.sourceClientId || data.sourceClientId === clientIdRef.current));
const pendingThreadId = pendingMessage?.threadId || current.activeThreadId;
if (keepPendingMessage && nextThreadId) {
await moveAgentUserMessage(pendingThreadId, nextThreadId, pendingMessage!.clientMessageId || pendingMessage!.itemId || pendingMessage!.id).catch(() => undefined);
}
applyWorkspaceChange(data);
if (!data.draftThread) void loadThreads(Boolean(data.emptyThread));
});
});
source.addEventListener("chat_message", (event) => {
const data = parseEventData<AgentChatEvent>(event);
if (!data?.message) return;
enqueueEvent(() => {
if (!isCurrentThreadEvent(data)) return;
if (!registerLiveAgentTurn(data, authoritativeHistoryTurnsRef.current, liveTurnKeysRef.current)) return;
const current = useAgentStore.getState();
const threadId = data.threadId || data.message!.threadId || current.activeThreadId;
const turnId = data.turnId ?? data.message!.turnId ?? "";
const clientMessageId = data.message!.clientMessageId || data.message!.itemId || data.message!.id;
if (current.activeThreadId !== threadId) return;
const pending = data.message!.role === "user"
? [...current.messages].reverse().find((item) => item.role === "user" && item.threadId === threadId && !item.turnId && (!clientMessageId || item.clientMessageId === clientMessageId))
: undefined;
const next = scopeChatItem({
...data.message!,
...(pending ? { clientMessageId: pending.clientMessageId, text: pending.text, historyText: pending.historyText, attachments: pending.attachments } : {}),
}, threadId, turnId);
const currentMessages = pending && turnId ? current.messages.filter((item) => item.id !== pending.id) : current.messages;
const messages = upsertAgentMessage(currentMessages, next);
setAgentState({ messages });
if (next.role === "user" && clientMessageId) void bindPendingAgentUserMessage(threadId, clientMessageId, turnId).catch(() => undefined);
if (next.role === "user" && !next.attachments?.length) {
void readAgentUserMessages(threadId).then((storedMessages) => {
const stored = storedMessages.find((item) => item.id === clientMessageId || Boolean(turnId && item.turnId === turnId));
const latest = useAgentStore.getState();
if (!stored || latest.activeThreadId !== threadId || !latest.messages.some((item) => item.id === next.id)) return;
setAgentState({ messages: upsertAgentMessage(latest.messages, { ...next, text: stored.text, historyText: stored.historyText, attachments: stored.attachments }) });
}).catch(() => undefined);
}
});
});
source.addEventListener("agent_log", (event) => {
if (!isCurrentConnection()) return;
const text = parseEventData<{ text?: unknown }>(event)?.text;
addEventLog("日志", text, text);
});
source.addEventListener("agent_error", (event) => {
const data = parseEventData<AgentEventPayload>(event);
if (!data) return;
enqueueEvent(() => {
if (!isCurrentThreadEvent(data)) return;
if (!registerLiveAgentTurn(data, authoritativeHistoryTurnsRef.current, liveTurnKeysRef.current)) return;
showAgentError(data.message, data, !data.replayed);
});
});
source.onerror = () => {
if (disposed || protocolRejected) return;
const wasConnected = connectedRef.current;
const silent = useAgentStore.getState().silentConnect && !wasConnected;
const text = wasConnected ? "本地 Agent 连接失败或已断开" : "连接失败,请检查地址和 token";
if (!errorLoggedRef.current || wasConnected) {
addEventLog(wasConnected ? "连接断开" : "连接失败", text);
if (!headless && !silent) message.error(text);
}
errorLoggedRef.current = true;
connectedRef.current = false;
pendingToolRef.current = null;
setAgentState({
activity: wasConnected ? "连接断开" : "连接失败",
connected: false,
waiting: false,
sending: false,
connectError: silent ? "" : text,
silentConnect: false,
pendingTool: null,
pendingApprovals: [],
});
if (!wasConnected) {
source.close();
setAgentState({ enabled: false });
}
};
return () => {
disposed = true;
source.close();
connectedRef.current = false;
loadThreadsSequenceRef.current += 1;
};
}, [applyWorkspaceChange, clientReady, enabled, endpoint, loadThreads, message, setAgentState, token]);
useEffect(() => {
if (connected) void loadThreads();
}, [connected, loadThreads]);
useEffect(() => {
if (!connected) return;
void fetchAgentJson<AgentModelsResponse>(endpoint, token, "/agent/codex/models").then(({ data = [] }) => {
const names = new Set<string>();
const models = data.flatMap((item) => {
const name = item.displayName || item.model;
const efforts = item.supportedReasoningEfforts.filter(({ reasoningEffort }) => AGENT_REASONING_EFFORTS.has(reasoningEffort));
if (item.model === "codex-auto-review" || names.has(name) || !efforts.length) return [];
names.add(name);
const defaultReasoningEffort = efforts.some((effort) => effort.reasoningEffort === item.defaultReasoningEffort) ? item.defaultReasoningEffort : efforts[0].reasoningEffort;
return [{ ...item, supportedReasoningEfforts: efforts, defaultReasoningEffort }];
});
if (!models.length) return;
const savedModel = useAgentStore.getState().model;
const current = models.find((item) => item.model === savedModel) || models.find((item) => item.isDefault) || models[0];
const savedEffort = useAgentStore.getState().reasoningEffort;
const efforts = current.supportedReasoningEfforts.map((item) => item.reasoningEffort);
const nextEffort = efforts.includes(savedEffort as AgentReasoningEffort) ? savedEffort as AgentReasoningEffort : current.defaultReasoningEffort || efforts[0];
localStorage.setItem("canvas-agent-model", current.model);
localStorage.setItem("canvas-agent-reasoning-effort", nextEffort);
setAgentState({ models, model: current.model, reasoningEffort: nextEffort });
}).catch((error) => addEventLog("读取模型列表失败", error));
}, [connected, endpoint, setAgentState, token]);
useEffect(() => {
if (!connected) return;
const activate = () => void activateAgentClient(endpoint, token, clientIdRef.current);
const activateVisible = () => {
if (document.visibilityState === "visible") activate();
};
window.addEventListener("focus", activate);
document.addEventListener("visibilitychange", activateVisible);
return () => {
window.removeEventListener("focus", activate);
document.removeEventListener("visibilitychange", activateVisible);
};
}, [connected, endpoint, token]);
const sendPrompt = async () => {
const text = prompt.trim();
const files = attachments;
const requestPrompt = promptWithAttachments(text, files);
const currentState = useAgentStore.getState();
if (!currentState.connected || !requestPrompt || currentState.sending || currentState.waiting || currentState.loadingThreads) return;
if (attachmentPayloadBytes(files) > MAX_ATTACHMENT_PAYLOAD_BYTES) {
addMessage({ role: "error", title: "图片过大", text: "图片附件超过 30MB,请删减后再发送。" });
return;
}
const messageId = createId();
const userText = text || `发送了 ${files.length} 张图片`;
loadThreadsSequenceRef.current += 1;
const currentBeforeSend = useAgentStore.getState();
const requestThreadId = currentBeforeSend.activeThreadId;
setAgentState({ prompt: "", attachments: [], activity: "发送中", sending: true, loadingThreads: false, activeTurnId: "", messages: currentBeforeSend.messages });
addMessage({ id: messageId, itemId: "synthetic:user", clientMessageId: messageId, threadId: requestThreadId, turnId: "", role: "user", text: userText, historyText: requestPrompt, attachments: files });
let threadId = requestThreadId;
try {
if (files.length) await savePendingAgentUserMessage({ id: messageId, role: "user", text: userText, historyText: requestPrompt, attachments: files });
const modelName = models.find((item) => item.model === model)?.displayName || model || "默认模型";
const effortName = reasoningEffort ? AGENT_REASONING_LABELS[reasoningEffort] : "默认强度";
addEventLog("发送任务", `${modelName} · ${effortName}${files.length ? ` · 附件 ${files.length}` : ""} · ${compactText(text) || "仅附件"}`);
const accepted = await fetchAgentJson<AgentTurnResponse>(endpoint, token, "/agent/codex/turn", {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify({
prompt: requestPrompt,
messageText: userText,
messageId,
clientId: clientIdRef.current,
threadId,
permissionMode,
model,
effort: reasoningEffort,
attachments: files.map(({ id, name, type, size, width, height, dataUrl }) => ({ id, name, type, size, width, height, dataUrl })),
}),
});
threadId = accepted.threadId || threadId;
if (!threadId) throw new Error("启动对话失败");
if (files.length) {
const latestMessage = useAgentStore.getState().messages.find((item) => item.clientMessageId === messageId);
const acceptedThreadId = latestMessage?.threadId || threadId;
await bindPendingAgentUserMessage(acceptedThreadId, messageId, latestMessage?.turnId || "").catch((error) => addEventLog("保存附件历史失败", error));
}
files.forEach((item) => {
URL.revokeObjectURL(item.url);
attachmentUrlsRef.current.delete(item.url);
});
} catch (error) {
if (files.length) await deletePendingAgentUserMessage(messageId).catch(() => undefined);
const text = error instanceof Error ? error.message : "发送失败";
const busy = text.includes("Codex 正在运行");
const state = useAgentStore.getState();
const removeFailedPending = (messages: AgentChatItem[]) => messages.filter((item) => item.clientMessageId !== messageId || Boolean(item.turnId));
threadMessagesRef.current.forEach((messages, cachedThreadId) => {
const next = removeFailedPending(messages);
if (next.length !== messages.length) threadMessagesRef.current.set(cachedThreadId, next);
});
const ownsCurrentThread = state.activeThreadId === (threadId || requestThreadId);
if (ownsCurrentThread) {
setAgentState({
activity: busy ? "Codex 正在运行" : "发送失败",
sending: false,
messages: removeFailedPending(state.messages),
...(state.prompt || state.attachments.length ? {} : { prompt, attachments: files }),
});
addMessage({ threadId: state.activeThreadId, turnId: "", role: "error", title: busy ? "任务仍在运行" : "发送失败", text });
} else {
setAgentState({ sending: false, messages: removeFailedPending(state.messages) });
}
addEventLog("发送失败", error);
}
};
const stopTurn = async () => {
if (!connected || (!sending && !waiting)) return;
setAgentState({ activity: "停止中" });
try {
await fetch(`${endpoint}/agent/codex/interrupt?token=${encodeURIComponent(token)}`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ threadId: useAgentStore.getState().activeThreadId || undefined }) });
addEventLog("停止任务", "已发送停止请求");
} catch {
setAgentState({ activity: "停止失败" });
}
};
const addAttachments = async (files: FileList | File[] | null) => {
if (!files) return;
const images = Array.from(files).filter((file) => file.type.startsWith("image/"));
const prev = useAgentStore.getState().attachments;
try {
const next = await Promise.all(
images.slice(0, Math.max(0, MAX_ATTACHMENTS - prev.length)).map(async (file) => {
const dataUrl = await readDataUrl(file);
const meta = await readImageMeta(dataUrl);
const url = URL.createObjectURL(file);
attachmentUrlsRef.current.add(url);
return { id: createId(), name: file.name, type: file.type, size: file.size, width: meta.width, height: meta.height, url, dataUrl };
}),
);
const merged = [...prev, ...next];
if (attachmentPayloadBytes(merged) > MAX_ATTACHMENT_PAYLOAD_BYTES) {
next.forEach((item) => {
URL.revokeObjectURL(item.url);
attachmentUrlsRef.current.delete(item.url);
});
addMessage({ role: "error", title: "图片过大", text: "图片附件最多约 30MB。" });
return;
}
if (next.length) setAgentState({ attachments: merged });
} catch (error) {
addMessage({ role: "error", title: "图片读取失败", text: error instanceof Error ? error.message : "图片读取失败" });
}
};
const removeAttachment = (id: string) => {
const removed = attachments.find((item) => item.id === id);
if (removed) {
URL.revokeObjectURL(removed.url);
attachmentUrlsRef.current.delete(removed.url);
}
setAgentState({ attachments: attachments.filter((item) => item.id !== id) });
};
const handleToolCall = async (endpoint: string, token: string, payload: AgentPendingToolCall) => {
if (confirmToolsRef.current && isCanvasWriteTool(payload.name)) {
if (pendingToolRef.current) {
await postToolResult(endpoint, token, clientIdRef.current, { requestId: payload.requestId, error: "仍有待确认的画布工具调用" });
return;
}
pendingToolRef.current = payload;
setAgentState({ pendingTool: payload });
addEventLog("等待确认", payload, payload);
return;
}
await runToolCall(endpoint, token, payload);
};
const runToolCall = async (endpoint: string, token: string, payload: AgentPendingToolCall) => {
if (isSiteTool(payload.name)) {
try {
addEventLog(toolName(payload.name), payload, payload);
const result = await runSiteTool(payload.name, payload.input || {}, navigate, { canvasSnapshot: canvasContextRef.current?.snapshot || null });
await postToolResult(endpoint, token, clientIdRef.current, { requestId: payload.requestId, result });
addEventLog(`${toolName(payload.name)}完成`, result, result);
} catch (error) {
const message = error instanceof Error ? error.message : "工具执行失败";
await postToolResult(endpoint, token, clientIdRef.current, { requestId: payload.requestId, error: message });
}
return;
}
try {
const input: { ops?: CanvasAgentOp[]; path?: string } = payload.input || {};
addEventLog(toolName(payload.name), payload, payload);
let result: unknown;
let appliedOps = input.ops || [];
if (payload.name === "site_navigate") {
const path = input.path || "/";
navigate(path);
result = { ok: true, path };
} else if (payload.name === "canvas_apply_ops") {
const context = canvasContextRef.current;
if (!context) throw new Error("当前不在画布页,请先用 site_navigate 打开画布");
result = context.applyOps(appliedOps);
void postState(endpoint, token, clientIdRef.current, result as CanvasAgentSnapshot);
} else if (payload.name === "canvas_create_attachment_nodes") {
const context = canvasContextRef.current;
if (!context) throw new Error("当前不在画布页,请先用 site_navigate 打开画布");
appliedOps = await attachmentNodeOps(endpoint, token, clientIdRef.current, payload.input?.nodes);
result = context.applyOps(appliedOps);
await postState(endpoint, token, clientIdRef.current, result as CanvasAgentSnapshot);
} else {
const snapshot = canvasContextRef.current?.snapshot;
if (!snapshot) throw new Error("当前不在画布页,请先用 site_navigate 打开画布");
result = snapshot;
}
await postToolResult(endpoint, token, clientIdRef.current, { requestId: payload.requestId, result });
addEventLog(`${toolName(payload.name)}完成`, result, result);
} catch (error) {
const message = error instanceof Error ? error.message : "画布操作失败";
await postToolResult(endpoint, token, clientIdRef.current, { requestId: payload.requestId, error: message });
}
};
const rejectPendingTool = async () => {
if (!pendingTool) return;
await postToolResult(endpoint, token, clientIdRef.current, { requestId: pendingTool.requestId, error: "用户取消了画布工具调用" });
pendingToolRef.current = null;
setAgentState({ pendingTool: null });
};
const approvePendingTool = async () => {
if (!pendingTool) return;
const tool = pendingTool;
pendingToolRef.current = null;
setAgentState({ pendingTool: null });
await runToolCall(endpoint, token, tool);
};
const decideApproval = async (approval: AgentPendingApproval, decision: "accept" | "acceptForSession" | "decline") => {
const current = useAgentStore.getState();
const pending = current.pendingApprovals.find((item) => item.requestId === approval.requestId);
if (!pending || pending.deciding) return;
setAgentState({ pendingApprovals: current.pendingApprovals.map((item) => item.requestId === approval.requestId ? { ...item, deciding: decision } : item), activity: "正在提交权限决定" });
try {
await postCodexApproval(endpoint, token, approval.requestId, decision);
const latest = useAgentStore.getState();
if (latest.pendingApprovals.some((item) => item.requestId === approval.requestId)) setAgentState({ activity: "等待 Codex 确认权限" });
} catch (error) {
const latest = useAgentStore.getState();
const expired = error instanceof Error && error.message.includes("审批请求已失效");
const resolved = expired || !latest.pendingApprovals.some((item) => item.requestId === approval.requestId);
const pendingApprovals = resolved
? latest.pendingApprovals.filter((item) => item.requestId !== approval.requestId)
: latest.pendingApprovals.map((item) => item.requestId === approval.requestId ? { ...item, deciding: undefined } : item);
setAgentState({ pendingApprovals, activity: approvalActivity(pendingApprovals, latest.waiting, latest.activity) });
if (resolved) return;
addEventLog("权限审批失败", error);
message.error(error instanceof Error ? error.message : "权限审批失败");
}
};
const changePermissionMode = (nextMode: AgentPermissionMode) => {
const apply = () => {
localStorage.setItem("canvas-agent-permission-mode", nextMode);
setAgentState({ permissionMode: nextMode });
};
if (nextMode !== "full") return apply();
modal.confirm({
title: "启用完全访问权限",
content: "Codex 将不受沙箱限制,可访问互联网及本机任意文件。请仅在信任当前任务时使用。",
okText: "启用完全访问",
okType: "danger",
cancelText: "取消",
onOk: apply,
});
};
const toggleAgentConnection = async ({ silent = false }: { silent?: boolean } = {}) => {
if (enabled) {
clearAgentSession({ enabled: false, connected: false, activity: "离线", connectError: "" });
return;
}
const urlToken = searchParams.get("agentToken") || "";
const urlEndpoint = searchParams.get("agentUrl") || "";
const discovered = urlToken ? null : await discoverAgentConfig(endpoint || DEFAULT_AGENT_URL);
const nextEndpoint = (urlEndpoint || discovered?.url || endpoint || DEFAULT_AGENT_URL).trim().replace(/\/$/, "");
const nextToken = (urlToken || token.trim() || discovered?.token || "").trim();
if (!nextEndpoint) {
const text = "请填写本地 Agent 地址";
if (!silent) {
setAgentState({ connectError: text });
if (!headless) message.warning(text);
}
return;
}
if (!nextToken) {
const text = "没有发现本地 Agent,请先在 Codex 使用插件或手动启动 Canvas Agent";
if (!silent) {
setAgentState({ connectError: text });
if (!headless) message.warning(text);
}
return;
}
try {
const parsed = new URL(nextEndpoint);
if (parsed.protocol !== "http:" && parsed.protocol !== "https:") throw new Error("invalid protocol");
} catch {
const text = "本地 Agent 地址格式不正确";
if (!silent) {
setAgentState({ connectError: text });
if (!headless) message.warning(text);
}
return;
}
errorLoggedRef.current = false;
setAgentState({ url: nextEndpoint, token: nextToken, enabled: true, connected: false, silentConnect: silent, activity: "连接中", connectError: "", activeTab: "setup" });
};
useEffect(() => {
if (urlAgentAutoConnect && confirmTools) setAgentState({ confirmTools: false });
}, [confirmTools, setAgentState, urlAgentAutoConnect]);
useEffect(() => {
if ((!autoConnect && !urlAgentAutoConnect) || autoConnectRef.current || enabled || connected) return;
autoConnectRef.current = true;
void toggleAgentConnection({ silent: true });
}, [autoConnect, connected, enabled, urlAgentAutoConnect]);
function clearAgentSession(patch: Parameters<typeof setAgentState>[0] = {}) {
loadThreadsSequenceRef.current += 1;
threadMessagesRef.current.clear();
authoritativeHistoryTurnsRef.current.clear();
liveTurnKeysRef.current.clear();
threadOperationRef.current = 0;
setAgentState({
messages: [],
tokenUsage: null,
threads: [],
activeThreadId: "",
activeTurnId: "",
workspacePath: "",
loadingThreads: false,
waiting: false,
sending: false,
pendingTool: null,
pendingApprovals: [],
bootstrapStatus: null,
mcpStartupStatuses: {},
...patch,
});
pendingToolRef.current = null;
}
const beginThreadOperation = () => {
const operation = ++threadOperationSequenceRef.current;
threadOperationRef.current = operation;
setAgentState({ loadingThreads: true });
return operation;
};
const finishThreadOperation = (operation: number) => {
if (threadOperationRef.current !== operation) return;
threadOperationRef.current = 0;
setAgentState({ loadingThreads: false });
};
const startNewThread = async () => {
const current = useAgentStore.getState();
if (!current.connected || current.sending || current.waiting || current.loadingThreads) return;
const operation = beginThreadOperation();
applyWorkspaceChange({ activeThreadId: "", emptyThread: true, draftThread: true, sourceClientId: clientIdRef.current });
setAgentState({ activeTab: "chat", activity: "正在新建对话", bootstrapStatus: { key: "codex:preparing", text: "正在初始化 Codex 对话", detail: "正在创建会话并启动画布工具服务", status: "running" }, mcpStartupStatuses: {} });
try {
const result = await fetchAgentJson<AgentWorkspaceResponse>(endpoint, token, "/agent/codex/threads/reset", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ clientId: clientIdRef.current, permissionMode }) });
if (threadOperationRef.current !== operation) return;
const latest = useAgentStore.getState();
if (latest.activeThreadId || latest.messages.length) applyWorkspaceChange({ activeThreadId: result.workspace?.activeThreadId || "", emptyThread: true, draftThread: true, sourceClientId: clientIdRef.current });
setAgentState({ activeTab: "chat", activity: "新对话" });
} catch (error) {
setAgentState({ bootstrapStatus: { key: "codex:prepare_failed", text: "Codex 对话初始化失败", detail: error instanceof Error ? error.message : "无法创建 Codex 会话", status: "error" } });
addEventLog("新建对话失败", error);
message.error(error instanceof Error ? error.message : "新建对话失败");
await loadThreads();
} finally {
finishThreadOperation(operation);
}
};
const resumeThread = async (threadId: string) => {
const current = useAgentStore.getState();
if (!current.connected || !threadId || current.sending || current.waiting || current.loadingThreads) return;
const operation = beginThreadOperation();
try {
await fetchAgentJson<AgentThreadResponse>(endpoint, token, `/agent/codex/threads/${encodeURIComponent(threadId)}/resume`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ permissionMode, clientId: clientIdRef.current }) });
await loadThreads();
if (useAgentStore.getState().activeThreadId === threadId) setAgentState({ activeTab: "chat", activity: "已恢复会话" });
} catch (error) {
addEventLog("恢复对话失败", error);
message.error(error instanceof Error ? error.message : "恢复对话失败");
await loadThreads();
} finally {
finishThreadOperation(operation);
}
};
const deleteThreads = async (threadIds: string[]) => {
if (!connected || !threadIds.length || sending || waiting || loadingThreads) return;
const operation = beginThreadOperation();
const deletedThreadIds: string[] = [];
try {
for (const threadId of new Set(threadIds)) {
await fetchAgentJson(endpoint, token, `/agent/codex/threads/${encodeURIComponent(threadId)}/delete`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ clientId: clientIdRef.current }) });
threadMessagesRef.current.delete(threadId);
deletedThreadIds.push(threadId);
}
void deleteAgentThreadMessages(deletedThreadIds).catch(() => undefined);
await loadThreads();
message.success(`已删除 ${deletedThreadIds.length} 条记录`);
} catch (error) {
if (deletedThreadIds.length) {
void deleteAgentThreadMessages(deletedThreadIds).catch(() => undefined);
}
await loadThreads();
addEventLog("删除对话失败", error);
message.error(error instanceof Error ? error.message : "删除对话失败");
} finally {
finishThreadOperation(operation);
}
};
const confirmDeleteThreads = (threadIds: string[]) => {
modal.confirm({
title: `删除 ${threadIds.length} 条对话记录`,
content: "删除后无法恢复,确定继续吗?",
okText: "删除",
okType: "danger",
cancelText: "取消",
onOk: () => deleteThreads(threadIds),
});
};
const addMessage = (item: Omit<AgentChatItem, "id"> & { id?: string }) => {
const text = normalizeText(item.text);
if (!text && !item.attachments?.length) return;
const current = useAgentStore.getState();
const itemId = item.itemId || item.id || createId();
const next = scopeChatItem({ ...item, id: item.id || itemId, itemId, text } as AgentChatItem, item.threadId ?? current.activeThreadId, item.turnId ?? current.activeTurnId);
setAgentState({ messages: upsertAgentMessage(current.messages, next) });
};
const addEventLog = (title: string, text: unknown, raw?: unknown) => {
const value = normalizeText(text) || title;
const last = useAgentStore.getState().eventLogs.at(-1);
if (last?.title === title && last.text === value) return;
pushEventLog({ id: `${Date.now()}-${Math.random()}`, time: dayjs().format("YYYY-MM-DD HH:mm:ss"), title, text: value, raw });
};
const upsertActivityMessage = (item: AgentChatItem) => {
setAgentState({ messages: upsertAgentMessage(useAgentStore.getState().messages, item) });
};
const appendActivityDelta = (event: AgentEventPayload) => {
const item = event.item;
if (!item?.id) return;
const text = stringText(item.text) || stringText(item.delta);
const isDelta = Boolean(stringText(item.delta));
if (!text) return;
if (item.type === "reasoning") {
const scoped = scopeEventChatItem(event, activityDeltaFallback(item, text), "synthetic:reasoning");
const current = useAgentStore.getState().messages.find((message) => message.id === scoped.id);
const activityItems = { ...(current?.activityItems || {}) };
const previous = activityItems[item.id] || "";
activityItems[item.id] = isDelta ? `${previous === activityPlaceholder("reasoning") ? "" : previous}${text}` : text;
upsertActivityMessage({ ...scoped, title: "思考摘要", text: reasoningActivityText(activityItems), activityItems, detail: activityDetail(current?.detail || scoped.detail, "reasoning", "inProgress") });
return;
}
const scoped = scopeEventChatItem(event, activityDeltaFallback(item, text), item.id);
const currentMessages = useAgentStore.getState().messages;
const index = currentMessages.findIndex((message) => message.id === scoped.id);
if (index < 0) {
if (!text.trim()) return;
upsertActivityMessage(scoped);
return;
}
const current = currentMessages[index];
if (item.type === "command_execution") {
const detail = activityDetail(current.detail, "command", "inProgress");
detail.output = isDelta ? `${stringText(detail.output)}${text}` : text;
setAgentState({ messages: currentMessages.map((message, itemIndex) => itemIndex === index ? { ...message, detail } : message) });
return;
}
const placeholder = activityPlaceholder(item.type);
if (!text.trim() && current.text === placeholder) return;
const nextText = isDelta ? `${current.text === placeholder ? "" : current.text}${text}` : mergeStreamText(current.text, text);
setAgentState({ messages: currentMessages.map((message, itemIndex) => itemIndex === index ? { ...message, text: nextText, detail: { ...activityDetail(message.detail, activityKind(item.type), "inProgress") } } : message) });
};
const upsertEventActivity = (event: AgentEventPayload, item: Omit<AgentChatItem, "id">) => {
const itemId = event.item?.id;
if (!itemId) return;
if (event.item?.type === "reasoning") {
const scoped = scopeEventChatItem(event, { ...item, id: "synthetic:reasoning" }, "synthetic:reasoning");
const current = useAgentStore.getState().messages.find((message) => message.id === scoped.id);
const activityItems = { ...(current?.activityItems || {}) };
const previous = activityItems[itemId] || "";
const incoming = normalizeText(item.text);
activityItems[itemId] = incoming === "已完成分析" && previous && previous !== activityPlaceholder("reasoning") ? previous : incoming;
upsertActivityMessage({ ...scoped, title: "思考摘要", text: reasoningActivityText(activityItems, incoming), activityItems });
return;
}
upsertActivityMessage(scopeEventChatItem(event, { ...item, id: itemId }, itemId));
};
const finishEmptyReasoningActivity = (event: AgentEventPayload) => {
const itemId = event.item?.id;
if (!itemId) return;
const scopedId = scopeEventChatItem(event, { id: "synthetic:reasoning", role: "tool", text: "" }, "synthetic:reasoning").id;
const currentMessages = useAgentStore.getState().messages;
const index = currentMessages.findIndex((message) => message.id === scopedId);
if (index < 0) return;
const current = currentMessages[index];
const activityItems = { ...(current.activityItems || {}) };
delete activityItems[itemId];
if (!Object.values(activityItems).some(isReasoningSummary)) {
setAgentState({ messages: currentMessages.filter((_, itemIndex) => itemIndex !== index) });
return;
}
setAgentState({ messages: currentMessages.map((message, itemIndex) => itemIndex === index ? { ...message, text: reasoningActivityText(activityItems), activityItems, detail: activityDetail(message.detail, "reasoning", "completed") } : message) });
};
const finishPlanActivity = (event: AgentEventPayload) => {
const id = scopeEventChatItem(event, { id: "synthetic:plan", role: "tool", text: "" }, "synthetic:plan").id;
const currentMessages = useAgentStore.getState().messages;
const index = currentMessages.findIndex((message) => message.id === id);
if (index < 0) return;
const current = currentMessages[index];
const detail = activityDetail(current.detail, "todo", turnPlanStatus(current.detail, event.status));
setAgentState({ messages: currentMessages.map((message, itemIndex) => itemIndex === index ? { ...message, detail } : message) });
};
const showAgentError = (value: unknown, event?: AgentEventPayload, log = true) => {
const error = agentErrorView(value);
const item = event
? scopeEventChatItem(event, { id: "synthetic:error", role: "error", title: error.title, text: error.text }, "synthetic:error")
: scopeChatItem({ id: createId(), role: "error", title: error.title, text: error.text }, useAgentStore.getState().activeThreadId, useAgentStore.getState().activeTurnId);
const state = useAgentStore.getState();
const current = state.messages.find((message) => message.id === item.id);
if (current && !normalizeText(value)) return;
upsertActivityMessage(item);
setAgentState({ activity: "处理失败", pendingApprovals: [] });
if (log) addEventLog("处理失败", error.text, value);
};
const handleAgentEvent = async (event: AgentEventPayload) => {
if (event.type === "usage.updated") setAgentState({ tokenUsage: eventUsage(event) });
const log = event.replayed ? null : formatAgentEventLog(event);
const activity = formatAgentActivity(event);
if (log) addEventLog(log.title, log.text);
if (event.type === "turn.started" && (event.turnId || event.turn_id)) {
const scope = eventScope(event);
const current = useAgentStore.getState();
if (!scope.threadId || !scope.turnId) return;
liveTurnKeysRef.current.add(`${scope.threadId}\0${scope.turnId}`);
setAgentState({ activeTurnId: scope.turnId, bootstrapStatus: null, mcpStartupStatuses: {}, messages: bindPendingTurnMessages(current.messages, scope.threadId, scope.turnId) });
}
if (event.type === "item.updated" && event.item?.type === "agent_message" && event.item.id) {
const delta = stringText(event.item.delta);
appendStreamText(event, delta || stringText(event.item.text), Boolean(delta));
return;
}
if (event.type === "item.updated" && event.item) {
appendActivityDelta(event);
return;
}
if (event.type === "plan.updated" && event.turn_id) {
const plan = formatAgentPlan(event);
if (plan) upsertActivityMessage(scopeEventChatItem(event, { ...plan, id: "synthetic:plan" }, "synthetic:plan"));
return;
}
if (event.type === "item.completed" && event.item?.type === "error") {
showAgentError(event.item.message, event, !event.replayed);
return;
}
if (event.type === "item.completed" && event.item?.type === "agent_message" && event.item.id) {
const scoped = scopeEventChatItem(event, { id: event.item.id, role: "assistant", title: "Codex", text: stringText(event.item.text) }, event.item.id);
const currentMessages = useAgentStore.getState().messages;
const index = currentMessages.findIndex((message) => message.id === scoped.id);
if (index >= 0) {
const text = stringText(event.item.text);
setAgentState({ messages: currentMessages.map((message, itemIndex) => itemIndex === index ? { ...message, text: text || message.text, streamId: undefined } : message) });
return;
}
addMessage(scoped);
return;
}
if (event.type === "item.completed" && event.item?.type === "reasoning" && !activity) {
finishEmptyReasoningActivity(event);
return;
}
if (event.type === "item.completed" && !activity && event.item?.id && event.item.type === "plan") {
const id = scopeEventChatItem(event, { id: event.item.id, role: "tool", text: "" }, event.item.id).id;
setAgentState({ messages: useAgentStore.getState().messages.filter((item) => item.id !== id) });
return;
}
if (!event.replayed && event.type === "item.completed" && event.item?.type === "image_generation" && event.item.id && event.sourceClientId === clientIdRef.current) {
const generated = await importGeneratedImages(endpoint, token, event.item);
if (generated.length) {
const context = canvasContextRef.current;
if (context) {
const right = Math.max(0, ...context.snapshot.nodes.map((node) => node.position.x + node.width)) + 80;
const ops = generated.map<CanvasAgentOp>((image, index) => {
const size = fitNodeSize(image.upload.width, image.upload.height);
return {
type: "add_node",
id: `image-${createId()}`,
nodeType: "image",
title: image.name,
position: { x: right + index * 40, y: index * 40 },
...size,
metadata: imageMetadata(image.upload),
};
});
const result = context.applyOps(ops);
void postState(endpoint, token, clientIdRef.current, result);
}
addEventLog("导入生成图片", context ? "已添加到发起任务的画布" : "图片已生成");
}
}
if (activity && event.item?.id) {
upsertEventActivity(event, activity);
return;
}
if (event.type === "turn.completed") {
const scope = eventScope(event);
if (scope.turnId) {
finishPlanActivity(event);
liveTurnKeysRef.current.add(`${scope.threadId}\0${scope.turnId}`);
}
const current = useAgentStore.getState();
setAgentState({
activeTurnId: current.activeTurnId === scope.turnId ? "" : current.activeTurnId,
messages: current.messages.map((message) => message.threadId === scope.threadId && message.turnId === scope.turnId && message.streamId ? { ...message, streamId: undefined } : message),
});
if (event.status === "failed") showAgentError(event.error?.message, event, !event.replayed);
}
const item = formatAgentEvent(event);
if (item) addMessage(scopeEventChatItem(event, { ...item, id: event.item?.id || createId() }, event.item?.id || createId()));
};
const appendStreamText = (event: AgentEventPayload, text: string, isDelta = false) => {
if (!text) return;
const itemId = event.item?.id;
if (!itemId) return;
const scoped = scopeEventChatItem(event, { id: itemId, role: "assistant", title: "Codex", text, streamId: itemId }, itemId);
const currentMessages = useAgentStore.getState().messages;
const index = currentMessages.findIndex((message) => message.id === scoped.id);
if (index < 0) {
pushMessage(scoped);
return;
}
setAgentState({ messages: currentMessages.map((message, itemIndex) => itemIndex === index ? { ...message, text: isDelta ? `${message.text}${text}` : mergeStreamText(message.text, text) } : message) });
};
const content = (
<>
<AgentPanelTabs
value={activeTab}
theme={theme}
leading={
<div className="flex items-center gap-2 pr-1">
<span className="grid size-8 place-items-center">
<Bot className="size-4" />
</span>
<div className="text-base font-semibold leading-5">Agent</div>
</div>
}
items={[
{ value: "setup", label: "连接", icon: <PlugZap className="size-3.5" /> },
{ value: "chat", label: "对话", icon: <MessageSquare className="size-3.5" /> },
{ value: "history", label: "历史", icon: <History className="size-3.5" />, count: threads.length },
{ value: "log", label: "日志", icon: <Terminal className="size-3.5" />, count: eventLogs.length },
]}
onChange={(activeTab) => {
setAgentState({ activeTab });
if (activeTab === "history") void loadThreads();
}}
right={
<>
<Button size="small" type="text" disabled={!connected || loadingThreads || sending || waiting} icon={<Plus className="size-3.5" />} onClick={startNewThread}>
新对话
</Button>
<Tooltip title="收起对话">
<Button type="text" shape="circle" className="!h-8 !w-8 !min-w-8" style={{ color: theme.node.muted }} icon={<PanelRightClose className="size-4" />} onClick={closePanel} />
</Tooltip>
</>
}
/>
{activeTab === "setup" ? (
<AgentConnectView
theme={theme}
url={url}
token={token}
enabled={enabled}
connected={connected}
activity={activity}
connectError={connectError}
onUrlChange={(url) => setAgentState({ url, connectError: "" })}
onTokenChange={(token) => setAgentState({ token, connectError: "" })}
onToggleEnabled={toggleAgentConnection}
/>
) : activeTab === "history" ? (
<AgentHistoryView
theme={theme}
threads={threads}
activeThreadId={activeThreadId}
workspacePath={workspacePath}
loading={loadingThreads}
busy={sending || waiting}
connected={connected}
onRefresh={() => void loadThreads()}
onNewThread={() => void startNewThread()}
onResumeThread={(threadId) => void resumeThread(threadId)}
onDeleteThreads={confirmDeleteThreads}
/>
) : activeTab === "log" ? (
<AgentLogView
logs={eventLogs}
theme={theme}
context={{ endpoint, connected, enabled, activity, waiting, sending, messages: messageCount, pendingTool: pendingTool?.name }}
onClear={clearEventLogs}
onCopied={(text) => message.success(text)}
onCopyBlocked={(text) => message.warning(text)}
/>
) : (
<>
<AgentChatTimeline theme={theme} pendingTool={pendingTool} pendingApprovals={pendingApprovals} sending={sending} waiting={waiting} onRejectTool={rejectPendingTool} onApproveTool={approvePendingTool} onApprovalDecision={decideApproval} />
<AgentTaskProgress theme={theme} busy={sending || waiting} />
{tokenUsage ? <AgentUsageBar usage={tokenUsage} theme={theme} /> : null}
<AgentChatComposer
prompt={prompt}
attachments={attachments.map(agentAttachmentToChatAttachment)}
disabled={!connected || agentInitializing}
sending={sending || waiting}
placeholder={agentInitializing ? "MCP 初始化中,完成后即可发送" : "询问 Codex,或让它操作网站/画布"}
theme={theme}
onPromptChange={(prompt) => setAgentState({ prompt })}
onSubmit={sendPrompt}
onStop={stopTurn}
onAddFiles={addAttachments}
onRemoveAttachment={removeAttachment}
confirmTools={confirmTools}
onConfirmToolsChange={(confirmTools) => setAgentState({ confirmTools })}
permissionMode={permissionMode}
onPermissionModeChange={changePermissionMode}
models={models}
model={model}
reasoningEffort={reasoningEffort}
onModelChange={(model) => {
const selected = models.find((item) => item.model === model);
if (!selected) return;
const effort = selected.defaultReasoningEffort || selected.supportedReasoningEfforts[0]?.reasoningEffort;
localStorage.setItem("canvas-agent-model", model);
if (effort) localStorage.setItem("canvas-agent-reasoning-effort", effort);
setAgentState({ model, ...(effort ? { reasoningEffort: effort } : {}) });
}}
onReasoningEffortChange={(reasoningEffort) => {
localStorage.setItem("canvas-agent-reasoning-effort", reasoningEffort);
setAgentState({ reasoningEffort });
}}
left={
attachments.length ? (
<span className="text-[11px]" style={{ color: theme.node.muted }}>
{formatBytes(attachmentPayloadBytes(attachments))} / 30MB
</span>
) : null
}
/>
</>
)}
</>
);
if (headless) return null;
return embedded ? content : null;
}
function acquireAgentClientId() {
const scope = globalThis as AgentClientGlobal;
scope.__infiniteCanvasAgentClientIdPromise ||= (async () => {
const storedClientId = readAgentClientId();
let clientId = storedClientId || randomId();
if (!navigator.locks) {
if (!storedClientId) saveAgentClientId(clientId);
return clientId;
}
while (true) {
const acquired = await new Promise<boolean>((resolve, reject) => {
void navigator.locks.request(`infinite-canvas-agent:${clientId}`, { ifAvailable: true }, async (lock) => {
if (!lock) return resolve(false);
resolve(true);
await new Promise<void>(() => undefined);
}).catch(reject);
});
if (acquired) {
saveAgentClientId(clientId);
return clientId;
}
clientId = randomId();
}
})().catch(() => {
const clientId = randomId();
saveAgentClientId(clientId);
return clientId;
});
return scope.__infiniteCanvasAgentClientIdPromise;
}
function readAgentClientId() {
try {
return sessionStorage.getItem("canvas-agent-client-id") || "";
} catch {
return "";
}
}
function saveAgentClientId(clientId: string) {
try {
sessionStorage.setItem("canvas-agent-client-id", clientId);
} catch {
// 内存身份仍可保证当前页面会话内的请求归属一致。
}
}
function eventScope(event: AgentEventPayload) {
return {
threadId: event.threadId || event.thread_id || "",
turnId: event.turnId || event.turn_id || "",
};
}
function scopeEventChatItem(event: AgentEventPayload, item: AgentChatItem, itemId: string) {
const scope = eventScope(event);
return scopeChatItem({ ...item, itemId }, scope.threadId, scope.turnId);
}
function approvalActivity(pendingApprovals: AgentPendingApproval[], waiting: boolean, fallback: string) {
if (pendingApprovals.length) return "等待权限确认";
return waiting ? "Codex 正在运行" : fallback;
}
async function attachmentNodeOps(endpoint: string, token: string, clientId: string, value: unknown): Promise<CanvasAgentOp[]> {
const nodes = Array.isArray(value) ? value : [];
if (!nodes.length) throw new Error("没有可添加的图片附件");
return await Promise.all(
nodes.map(async (value) => {
const item = value as { id?: unknown; attachmentId?: unknown; title?: unknown; position?: unknown };
const id = String(item.id || "");
const attachmentId = String(item.attachmentId || "");
if (!id || !attachmentId) throw new Error("图片附件节点参数无效");
const res = await fetch(`${endpoint}/agent/attachments/${encodeURIComponent(attachmentId)}?token=${encodeURIComponent(token)}&clientId=${encodeURIComponent(clientId)}`);
if (!res.ok) {
const body = (await res.json().catch(() => null)) as { error?: string } | null;
throw new Error(body?.error || "读取图片附件失败");
}
const image = await uploadImage(await res.blob());
const size = fitNodeSize(image.width, image.height);
const position = item.position && typeof item.position === "object" ? (item.position as { x?: unknown; y?: unknown }) : {};
return {
type: "add_node" as const,
id,
nodeType: "image" as const,
title: String(item.title || "参考图"),
position: { x: Number(position.x) || 0, y: Number(position.y) || 0 },
width: size.width,
height: size.height,
metadata: imageMetadata(image),
};
}),
);
}
function createId() {
return randomId();
}
function clamp(value: number, min: number, max: number) {
return Math.min(max, Math.max(min, value));
}
async function importGeneratedImages(endpoint: string, token: string, item: AgentEventItem) {
const sources = Array.from(generatedImageSources(item));
return await Promise.all(
sources.map(async (source, index) => {
const response = source.startsWith("data:image/")
? await fetch(source)
: await fetch(`${endpoint}/agent/local-image?token=${encodeURIComponent(token)}`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ path: source }) });
if (!response.ok) throw new Error("读取 Codex 生成图片失败");
const blob = await response.blob();
const upload = await uploadImage(blob);
const dataUrl = await readDataUrl(blob);
const name = source.startsWith("/") ? source.split("/").at(-1) || `生成图片 ${index + 1}` : `生成图片 ${index + 1}`;
return { upload, name, attachment: { id: createId(), name, type: blob.type || upload.mimeType, size: blob.size, width: upload.width, height: upload.height, url: upload.url, dataUrl } };
}),
);
}
function generatedImageSources(value: unknown, result = new Set<string>()) {
if (typeof value === "string") {
if (value.startsWith("data:image/") || (/^\/.+\.(?:avif|gif|jpe?g|png|webp)$/i.test(value) && !value.includes("\n"))) result.add(value);
return result;
}
if (Array.isArray(value)) value.forEach((item) => generatedImageSources(item, result));
else if (value && typeof value === "object") Object.values(value).forEach((item) => generatedImageSources(item, result));
return result;
}
function readDataUrl(file: Blob) {
return new Promise<string>((resolve, reject) => {
const reader = new FileReader();
reader.onload = () => resolve(String(reader.result || ""));
reader.onerror = () => reject(reader.error || new Error("读取图片失败"));
reader.readAsDataURL(file);
});
}
function delay(ms: number) {
return new Promise<void>((resolve) => setTimeout(resolve, ms));
}