| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369 |
- import { shell } from 'electron';
- import { logger } from 'ee-core/log';
- import { registerArtifactRoot } from '../service/codex/codexArtifactService';
- import { getCodexHome, getSkillsDir, isInsideDir } from '../service/codex/codexHome';
- import { probeCodexBinary } from '../service/codex/codexLocator';
- import {
- applyProvider as applyProviderToRuntime,
- baseUrlHint,
- clearProvider as clearRuntimeProvider,
- describeDroppedKeys,
- probeEndpoint,
- } from '../service/codex/providerService';
- import {
- defaultWorkspaceDir,
- getAppliedProvider,
- getCodex,
- setAppliedProvider,
- tailDiagnostics,
- toStatusResult,
- } from '../service/codex/index';
- import {
- asMessage,
- fail,
- ok,
- type ApprovalAnswers,
- type ApprovalDecision,
- type ApplyProviderInput,
- type CodexPingResult,
- type CodexStatusResult,
- type Rpc,
- } from '../service/codex/types';
- import type { ModelOption } from '../service/codex/codexRuntime';
- import type { McpServerInput, McpServerSetting } from '../service/codex/mcpService';
- import type { SkillListItem, SkillReadResult } from '../service/codex/skillService';
- import type { AgentEvent } from '../service/codex/types';
- const APPROVAL_DECISIONS: readonly ApprovalDecision[] = ['accept', 'acceptForSession', 'decline', 'cancel'];
- /** 统一的异常包装:ee-core 的 ipcMain.handle 不 catch,抛错会被 Electron 包成 "Error occurred in handler for ..." */
- async function guard<T>(label: string, run: () => Promise<T>): Promise<Rpc<T>> {
- try {
- return ok(await run());
- } catch (error) {
- const message = asMessage(error);
- logger.error(`[codexCtl] ${label} failed:`, message);
- return fail(message);
- }
- }
- function requireString(value: unknown, label: string): string {
- if (typeof value !== 'string' || !value.trim()) throw new Error(`${label} 不能为空`);
- return value.trim();
- }
- /** Codex 的「线程不在本进程里」只以这句错误文本出现,是识别失效线程的唯一依据 */
- function isThreadMissingError(error: unknown): boolean {
- return /thread not found/iu.test(error instanceof Error ? error.message : String(error));
- }
- function toBytes(value: unknown): Uint8Array {
- if (value instanceof Uint8Array) return value;
- if (value instanceof ArrayBuffer) return new Uint8Array(value);
- if (Array.isArray(value)) return new Uint8Array(value as number[]);
- if (value && typeof value === 'object' && 'data' in value) {
- return toBytes((value as { data: unknown }).data);
- }
- throw new Error('文件内容必须是二进制数据');
- }
- class CodexCtl {
- /** 只探测二进制与 CODEX_HOME,不 spawn 子进程 */
- async ping(): Promise<Rpc<CodexPingResult>> {
- return guard('ping', async () => {
- // getCodex() 会顺带建好 CODEX_HOME 与 skills 目录
- await getCodex();
- const probe = await probeCodexBinary();
- return { ...probe, codexHome: getCodexHome() };
- });
- }
- async status(): Promise<Rpc<CodexStatusResult>> {
- return guard('status', async () => {
- const { runtime } = await getCodex();
- await getAppliedProvider();
- return toStatusResult(runtime.status);
- });
- }
- async start(): Promise<Rpc<CodexStatusResult>> {
- return guard('start', async () => {
- const { runtime } = await getCodex();
- await getAppliedProvider();
- return toStatusResult(await runtime.start());
- });
- }
- async stop(): Promise<Rpc<CodexStatusResult>> {
- return guard('stop', async () => {
- const { runtime } = await getCodex();
- await runtime.stop();
- return toStatusResult(runtime.status);
- });
- }
- async restart(): Promise<Rpc<CodexStatusResult>> {
- return guard('restart', async () => {
- const { runtime } = await getCodex();
- await runtime.stop();
- await getAppliedProvider();
- return toStatusResult(await runtime.start());
- });
- }
- async listModels(): Promise<Rpc<ModelOption[]>> {
- return guard('listModels', async () => {
- const { runtime } = await getCodex();
- return runtime.listModels();
- });
- }
- async providerCapabilities(): Promise<Rpc<unknown>> {
- return guard('providerCapabilities', async () => {
- const { runtime } = await getCodex();
- return runtime.readModelProviderCapabilities();
- });
- }
- /** 只做端点探测,不改运行时;页面在「应用」前给用户预览兼容性 */
- async probeProvider(params: {
- baseUrl?: string;
- apiKey?: string | null;
- modelId?: string | null;
- }): Promise<Rpc<unknown>> {
- return guard('probeProvider', async () => {
- const baseUrl = requireString(params?.baseUrl, 'baseUrl');
- const result = await probeEndpoint({
- baseUrl,
- modelId: params?.modelId ?? null,
- apiKey: params?.apiKey ?? null,
- });
- return { ...result, hint: baseUrlHint(baseUrl) };
- });
- }
- /** 应用后端「模型管理」里的一条记录;会重启子进程(api key 走环境变量) */
- async applyProvider(params: ApplyProviderInput & { modelRecordId?: string | number | null }): Promise<Rpc<unknown>> {
- return guard('applyProvider', async () => {
- if (!params || typeof params !== 'object') throw new Error('缺少模型配置');
- const { runtime } = await getCodex();
- const applied = await applyProviderToRuntime(runtime, params);
- setAppliedProvider(applied);
- return {
- applied,
- hint: baseUrlHint(applied.baseUrl),
- droppedConfigKeys: describeDroppedKeys(params),
- status: toStatusResult(runtime.status),
- };
- });
- }
- async clearProvider(): Promise<Rpc<CodexStatusResult>> {
- return guard('clearProvider', async () => {
- const { runtime } = await getCodex();
- await clearRuntimeProvider(runtime);
- setAppliedProvider(null);
- return toStatusResult(runtime.status);
- });
- }
- async mcpList(): Promise<Rpc<McpServerSetting[]>> {
- return guard('mcpList', async () => (await getCodex()).mcp.list());
- }
- async mcpSave(params: McpServerInput): Promise<Rpc<McpServerSetting[]>> {
- return guard('mcpSave', async () => (await getCodex()).mcp.save(params));
- }
- async mcpRemove(params: { id?: string }): Promise<Rpc<McpServerSetting[]>> {
- return guard('mcpRemove', async () =>
- (await getCodex()).mcp.remove(requireString(params?.id, 'id')),
- );
- }
- async skillList(params?: { forceReload?: boolean }): Promise<Rpc<SkillListItem[]>> {
- return guard('skillList', async () =>
- (await getCodex()).skills.list({ forceReload: params?.forceReload === true }),
- );
- }
- async skillRead(params: { path?: string; name?: string }): Promise<Rpc<SkillReadResult>> {
- return guard('skillRead', async () => (await getCodex()).skills.read(params ?? {}));
- }
- async skillInstallFolder(params: {
- srcPath?: string;
- name?: string | null;
- overwrite?: boolean;
- }): Promise<Rpc<SkillListItem[]>> {
- return guard('skillInstallFolder', async () =>
- (await getCodex()).skills.installFromFolder(requireString(params?.srcPath, 'srcPath'), {
- name: params?.name ?? null,
- overwrite: params?.overwrite === true,
- }),
- );
- }
- async skillInstallZip(params: {
- data?: unknown;
- name?: string | null;
- overwrite?: boolean;
- }): Promise<Rpc<SkillListItem[]>> {
- return guard('skillInstallZip', async () =>
- (await getCodex()).skills.installFromZip(toBytes(params?.data), {
- name: params?.name ?? null,
- overwrite: params?.overwrite === true,
- }),
- );
- }
- async skillRemove(params: { path?: string; name?: string }): Promise<Rpc<SkillListItem[]>> {
- return guard('skillRemove', async () => (await getCodex()).skills.remove(params ?? {}));
- }
- async skillSetEnabled(params: {
- path?: string;
- name?: string;
- enabled?: boolean;
- }): Promise<Rpc<{ effectiveEnabled: boolean }>> {
- return guard('skillSetEnabled', async () => {
- if (typeof params?.enabled !== 'boolean') throw new Error('enabled 必须是布尔值');
- const effectiveEnabled = await (await getCodex()).skills.setEnabled(params, params.enabled);
- return { effectiveEnabled };
- });
- }
- /** 只允许打开用户 skills 目录或其中的某个 skill */
- async skillOpenFolder(params?: { path?: string }): Promise<Rpc<{ path: string }>> {
- return guard('skillOpenFolder', async () => {
- const skillsDir = getSkillsDir();
- const target = params?.path ? params.path : skillsDir;
- if (!isInsideDir(skillsDir, target) && target !== skillsDir) {
- throw new Error('只能打开 skills 目录内的路径');
- }
- const result = await shell.openPath(target);
- if (result) throw new Error(result);
- return { path: target };
- });
- }
- /** 会话默认工作目录(渲染进程也需要知道真实值,产物栏要与 cwd 对齐) */
- async defaultWorkspace(): Promise<Rpc<{ cwd: string }>> {
- return guard('defaultWorkspace', async () => ({ cwd: await defaultWorkspaceDir() }));
- }
- async threadStart(params?: {
- cwd?: string | null;
- model?: string | null;
- approvalPolicy?: 'untrusted' | 'on-request' | 'never';
- sandbox?: 'read-only' | 'workspace-write' | 'danger-full-access';
- }): Promise<Rpc<{ threadId: string; cwd: string }>> {
- return guard('threadStart', async () => {
- const { runtime } = await getCodex();
- const cwd = params?.cwd?.trim() || (await defaultWorkspaceDir());
- const threadId = await runtime.startThread({
- cwd,
- model: params?.model ?? null,
- approvalPolicy: params?.approvalPolicy ?? 'never',
- sandbox: params?.sandbox ?? 'workspace-write',
- });
- // 登记产物白名单:产物读写只允许发生在会话工作目录内
- registerArtifactRoot(cwd);
- return { threadId, cwd };
- });
- }
- /** 跑一整轮并等结果;事件流同时通过 codex/event 推给渲染进程 */
- async turnRun(params: {
- threadId?: string;
- input?: string;
- cwd?: string | null;
- model?: string | null;
- approvalPolicy?: 'untrusted' | 'on-request' | 'never';
- timeoutMs?: number;
- }): Promise<Rpc<{ turnId: string; status: string; text: string }>> {
- return guard('turnRun', async () => {
- const { runtime } = await getCodex();
- const threadId = requireString(params?.threadId, 'threadId');
- const cwd = params?.cwd?.trim() || undefined;
- if (cwd) registerArtifactRoot(cwd);
- const options = {
- threadId,
- prompt: requireString(params?.input, 'input'),
- cwd,
- model: params?.model ?? null,
- approvalPolicy: params?.approvalPolicy ?? ('never' as const),
- timeoutMs: typeof params?.timeoutMs === 'number' ? params.timeoutMs : undefined,
- };
- try {
- const result = await runtime.runTurn(options);
- return { turnId: result.turnId, status: result.status, text: result.text };
- } catch (error: unknown) {
- // 线程只存在于当前 app-server 进程里:换模型或重启应用后旧 threadId 必然认不出,
- // 先按 Codex 落盘的历史把它救回来(thread/resume),再重发这一轮
- if (!isThreadMissingError(error)) throw error;
- try {
- await runtime.resumeThread(threadId, {
- cwd: cwd || (await defaultWorkspaceDir()),
- model: options.model,
- approvalPolicy: options.approvalPolicy,
- });
- } catch {
- // 从没成功跑完过一轮的线程 Codex 根本没落盘,救不回来
- throw new Error('该会话的 Codex 线程已失效(中途换过模型或重启过应用),请新建会话继续');
- }
- const retried = await runtime.runTurn(options);
- return { turnId: retried.turnId, status: retried.status, text: retried.text };
- }
- });
- }
- async turnInterrupt(params: { threadId?: string; turnId?: string }): Promise<Rpc<Record<string, never>>> {
- return guard('turnInterrupt', async () => {
- const { runtime } = await getCodex();
- await runtime.interruptTurn(requireString(params?.threadId, 'threadId'), requireString(params?.turnId, 'turnId'));
- return {};
- });
- }
- async threadUnsubscribe(params: { threadId?: string }): Promise<Rpc<Record<string, never>>> {
- return guard('threadUnsubscribe', async () => {
- const { runtime } = await getCodex();
- await runtime.unsubscribeThread(requireString(params?.threadId, 'threadId'));
- return {};
- });
- }
- async approvalResolve(params: {
- token?: string;
- decision?: ApprovalDecision;
- answers?: ApprovalAnswers;
- }): Promise<Rpc<Record<string, never>>> {
- return guard('approvalResolve', async () => {
- const { broker } = await getCodex();
- const token = requireString(params?.token, 'token');
- const decision = params?.decision;
- if (!decision || !APPROVAL_DECISIONS.includes(decision)) throw new Error('decision 不合法');
- broker.resolve(token, decision, params?.answers);
- return {};
- });
- }
- /** 读回某个 thread 已落盘的事件(页面刷新后恢复现场) */
- async eventsRead(params: { threadId?: string; limit?: number }): Promise<Rpc<AgentEvent[]>> {
- return guard('eventsRead', async () => {
- const { eventLog } = await getCodex();
- const threadId = requireString(params?.threadId, 'threadId');
- return eventLog.read(threadId, typeof params?.limit === 'number' ? params.limit : undefined);
- });
- }
- async logsTail(params?: { limit?: number }): Promise<Rpc<Array<{ ts: string; line: string }>>> {
- return guard('logsTail', async () =>
- tailDiagnostics(typeof params?.limit === 'number' ? params.limit : undefined),
- );
- }
- }
- CodexCtl.toString = () => 'CodexCtl';
- export default CodexCtl;
|