diff --git a/.project-docs/30-worklog/tasks/20260910-personal-agent-platform-d15b2559.md b/.project-docs/30-worklog/tasks/20260910-personal-agent-platform-d15b2559.md new file mode 100644 index 0000000..4868066 --- /dev/null +++ b/.project-docs/30-worklog/tasks/20260910-personal-agent-platform-d15b2559.md @@ -0,0 +1,47 @@ +# Task: Complete personal cloud Agent platform across WS Yuxi and MakeLore + +## Identity + +- Task ID: 20260910-personal-agent-platform-d15b2559 +- Mode: Feature +- Branch: codex/20260910-personal-agent-platform-d15b2559-personal-agent-platform +- Worktree: D:\Datas\OthersProjects\.codex-worktrees\makelore\20260910-personal-agent-platform-d15b2559 +- Base commit: 46fbae9dec70ee1bbfae671d458a81fd1d8a1c28 +- Owner: codex +- Status: Ready for Integration + +## Scope + +- Complete the native cloud Agent workspace: configuration/preview, conversations/files/approval, publication/user sharing/application API, schedules/activity/cost; all cloud traffic and credentials remain Main-owned. +- User explicitly requested all remaining first-release modules, not another draft-only slice. Build on the reviewed entry commit recorded above; verify all accepted personal/share/API/scheduled/goal-tool paths before claiming completion. + +## Intent And Constraints + +- Bundled Concurrent Task Gate and Planning Gate passed. All operations are rooted in this owned isolated feature worktree. Peer records were read, with only the completed entry task added since the previous assessment; unrelated and defective peers remain read-only. +- Personal creators only; no organization/department/team wallet or Xiaozhi integration. WS owns identity/module/model eligibility and final Token Points; Yuxi owns cloud Request/Run/scheduler/resources. Existing WS Code/Canvas Gateway continues its own scope; this is the user-approved cloud-module extension recorded in the three-end plan. +- Creator, caller and payer remain distinct. Shared/API callers keep private conversations/files; every real model call charges the creator via existing WS accounting, never falls back to charging the caller. Application subjects are not fake personal accounts. +- Reuse Yuxi configuration, repositories, durable FIFO/Run/worker/approval, workspace and scheduling services; explicit empty capability selection disables it. Publish immutable configurations and freeze accepted requests, including resume/subcalls. Main owns cloud credentials/network, Renderer consumes safe Host API projections. +- Plan: verify and close model delegation/application scope contracts; implement server persistence/execution/finance and real transaction/worker tests; implement all native client surfaces and relevant Electron flows; run repository gates and document honest validation limits. No production data operations, deployment, push or additional subagents are implied. + +## Outcome + +- Completed native Agents workspace: My/Received/Activity, configuration and preview, formal conversations, approval/questions, attachments and downloadable artifacts, personal KB documents/indexing, publication/share/application credentials, creator costs and scheduled goals/history. +- Main owns exchanged credentials, actual cloud transport/SSE, account cancellation, local file dialogs/writes and desktop notifications; Renderer uses typed Host operations. Unsaved/unsent input and uncertain requests are preserved through explicit leave/retry flows; niancode Agent links route through login and server authorization. +- Implementation and verification are complete. The user explicitly authorized one fresh independent read-only Reviewer and commits after fixes. The fresh Reviewer reviewed the full three-repository scope and returned final PASS after confirming all findings were fixed. No unresolved review finding remains; no production deployment, push or merge was performed. + +## Verification + +- Passed: 32 tests across Main, draft/UI workflows and deep links. Includes stable retry intent, native batched SSE, approval restore, schedule enable intent, cost-service partial failure, actual temporary file bytes through mocked native dialogs, account switch cancellation and notification deduplication. +- Review fixes verified: new conversations appear immediately in history; unsent preview messages survive configuration edits; conversation changes protect approval answers and temporary attachments; cancelling file selection creates no empty conversation. The final attachment case verifies that only confirmation creates the conversation and associates the file. Renderer typecheck, targeted ESLint and Vite build passed again after this last fix. +- Passed: Renderer TypeScript; Vite build for Renderer/Main/preload/utility; real Electron Playwright scenario creates/saves/publishes an Agent, shows formal Markdown chat and creates an enabled schedule. The Host backend is explicitly stubbed in Electron coverage. +- Passed: targeted ESLint with no errors (2 non-blocking CloudChat warnings). Main TypeScript has the same 61 pre-existing errors, none in cloud Agents; this does not represent a clean repository-wide Main typecheck. +- Electron screenshots were inspected. No packaged installation, OS protocol registration or live three-service integration was executed. + +## Follow-ups + +- Hand off the reviewed result as committed feature branches for a separately authorized integration/release. Keep the worktrees and branches; no push, merge, deployment or further Reviewer is implied. +- Release-environment follow-up: actual one-api/paid-model reconciliation, provisioner sandbox execution, installable client packaging and production migration/rollback exercise. Current local validation does not claim these deployment outcomes. + +## Promotion Candidates + +- Target: integrated state and owning Agent/auth/billing/runtime architecture records. Proposal: complete the user-approved personal cloud platform using existing Yuxi capability/Request/Run/scheduling and WS accounting owners; document the cloud module's scope alongside the existing WS Gateway. Evidence: this task's final source and meaningful HTTP/PG/worker/Electron validation. Future impact: future modules consume one runtime and one ledger. Semantic conflicts: earlier planning fields and draft-only scope are superseded by implementation; existing unrelated product contracts remain. Human confirmation: product direction is already approved; canonical promotion remains a serialized integration task. diff --git a/README.md b/README.md index 48ca473..1cad1f8 100644 --- a/README.md +++ b/README.md @@ -5,11 +5,11 @@ Makelore 是一个面向软件、视觉创作、智能机器人与个人云智能体的 AI 桌面工作台。当前版本为 `2.0.0`,源码提供四个模块入口;云智能体需配套配置 WS/Yuxi 服务。模块入口页采用统一的横向卡片视觉,工作区左上角入口点击后返回模块入口页: - `Makelore Code|AI 编程`:管理本地项目、项目智能体、对话、文件上下文、代码变更和运行时。 -- `Makelore Agents|AI 智能体`:使用现有账号创建私人云草稿,编辑名称、用途和角色指令。Main 持有 WS/Yuxi 云凭据,Renderer 只通过 Host API 获取草稿;保存检查修订号,冲突保留本地输入,离开时提示未保存修改。当前没有试聊、发布、分享、API 调用和自动任务;保存不扣词元点数。 +- `Makelore Agents|AI 智能体`:配置个人云智能体的模型、角色、工具、知识库、MCP、Skill、子智能体与执行限制;保存草稿并试聊,发布后正式对话、按账号分享或创建应用 API Key。自动任务支持目标、时间规则、审批恢复与运行历史;活动页汇总未读结果。对话支持流式回复、断线恢复、附件和产物下载,个人资料可上传并建立知识索引。所有模型调用由创建者个人词元点数支付,分享/API 调用者保有自己的内容空间;保存配置不扣点。Main 管理云会话和本机文件选择,Renderer 通过 Host API 操作。配套服务接入见 [Yuxi MakeLore 说明](https://xerrors.github.io/Yuxi/advanced/makelore-agents.html)。 - `Makelore Canvas|AI 绘画`:每个设计项目(Workspace)维护一份从创建起就存在的 Living Form。左侧项目栏负责新建、切换和管理 Workspace,并在桌面设计模式下以 256px 宽度常驻展开;中央沿用 AI 编程的安静对话画布、自然消息流和底部悬浮输入器,AI 整理出的制作方案作为对话内的轻量可编辑稿持续更新;桌面端右侧同为 256px 的全高历史作品栏集中展示当前项目的制作记录与生成结果。紧凑窗口通过左侧抽屉访问项目列表,历史记录保留在时间线中。参考图从本地上传后以 `@图片N` 绑定,具体用法只写在创作提示词中。 - `Makelore Robot|AI 机器`:管理机器人智能体、设备激活绑定、智能体配置与设备分配;机器人工作台的智能体位于 Robot 全局侧栏,选中后在内容区先查看绑定设备、再查看基础设置,当前智能体通过 URL 参数保持可分享选择;绑定设备时默认先选择“引导配网”或“已有激活码”。在 Windows 与 macOS 的引导路径中,Makelore 可在弹窗内扫描并连接附近开放的 `Xiaozhi-*` 配网热点,失败时仍可通过系统 Wi-Fi 手动连接;后续继续复用机器人现有热点配网页面,不修改固件,也不由 Makelore 接收 Wi-Fi 密码。 -应用启动默认进入 AI 模块入口选择页。入口页可在未登录状态浏览;未登录用户点击已开通模块时进入客户端原生登录页,可使用账号密码或手机号短信验证码登录。密码登录可选“记住密码”:正式安装包仅由 Electron Main 使用系统受保护凭据存储加密保存和回填账号密码,不写入 Renderer 持久状态,未打包开发版或系统安全存储不可用时禁用该选项。登录请求由 Renderer 经 Host API 交给 Electron Main,再由 Main 调用 Works Square;成功后回到入口选择页。已登录时,Electron Main 会从 Works Square `/api/auth/me` 读取当前账号,只向 Renderer 投影用户名、账号/租户/部门标识、权限名列表与三个模块布尔开关,不透传上游资料或凭据。工作区门禁同时要求有效 Token 和完整用户身份;旧状态缺失身份时会先尝试从 Main 恢复,仍无法确认则清除残留会话并返回登录页。被管理员关闭的模块会在入口页置灰且无法点击,直接访问其工作区路径也会返回入口页。旧服务端未返回策略或缺少单项字段时默认开放;这个客户端门禁不替代服务端 API 授权。 +应用启动默认进入 AI 模块入口选择页。入口页可在未登录状态浏览;未登录用户点击已开通模块时进入客户端原生登录页,可使用账号密码或手机号短信验证码登录。密码登录可选“记住密码”:正式安装包仅由 Electron Main 使用系统受保护凭据存储加密保存和回填账号密码,不写入 Renderer 持久状态,未打包开发版或系统安全存储不可用时禁用该选项。登录请求由 Renderer 经 Host API 交给 Electron Main,再由 Main 调用 Works Square;成功后回到入口选择页。已登录时,Electron Main 会从 Works Square `/api/auth/me` 读取当前账号,只向 Renderer 投影用户名、账号/租户/部门标识、权限名列表与四个模块布尔开关,不透传上游资料或凭据。工作区门禁同时要求有效 Token 和完整用户身份;旧状态缺失身份时会先尝试从 Main 恢复,仍无法确认则清除残留会话并返回登录页。被管理员关闭的模块会在入口页置灰且无法点击,直接访问其工作区路径也会返回入口页。旧服务端未返回策略或缺少单项字段时默认开放;这个客户端门禁不替代服务端 API 授权。 作品广场、素材广场、独立发布上传和云部署页面不属于 Makelore 2.0 工作台。新建 Code 项目只要求选择目录:Main 自动生成内部项目 ID,并以内部 `interactive_ai_app` 类型创建 `.makelore/project.json` 与 `knowledge/`,不再让用户选择或查看项目身份、项目类型和模板;缺少项目 ID 的旧项目在读取时由 Main 自动补全。创建成功后直接进入对话工作区,未创建智能体时只显示可选的设置入口,不再用初始化门禁遮挡工作区。已有 `custom` 项目继续受支持;历史 `mini_game` / `mini_program` 配置在读取时归一为交互式 AI 应用,但不会因读取被改写。用户获取并为项目启用官方 bundled `makelore.project-scaffold` 插件后,每个父智能体都可按需明确调用 `makelore-project-scaffold` Skill,无需伙伴分配;它以不覆盖既有路径的方式生成固定六文件 Vite 起步工程,不是创建前置条件,也不安装依赖、不联网、不构建、不上传或提审。交互式 AI 应用的项目配置底部提供“一键提交审核”;Main 自动预检、安全打包并提交,构建通过后进入运营审核,审核通过即直接发布。首次创建必须选择 PNG、JPEG 或 WebP 项目封面,并通过 Main-owned multipart 原子接口同时保存资料与封面;已有 draft/published 只提交新版本并沿用平台现有资料与封面。项目成果预览 `/deliverables` 继续保留。 diff --git a/electron/api/routes/cloud-agents.ts b/electron/api/routes/cloud-agents.ts index 67c5da3..d227d6f 100644 --- a/electron/api/routes/cloud-agents.ts +++ b/electron/api/routes/cloud-agents.ts @@ -1,7 +1,7 @@ import type { IncomingMessage, ServerResponse } from 'node:http'; import { CLOUD_AGENTS_PATH } from '../../../shared/cloud-agents'; import { CloudAgentsError, CloudAgentsModule } from '../../services/cloud-agents'; -import { parseJsonBody, sendJson } from '../route-utils'; +import { parseJsonBody, sendJson, flushStreamingHeaders, writeStreamingChunk } from '../route-utils'; const cloudAgents = new CloudAgentsModule(); @@ -11,7 +11,39 @@ export async function handleCloudAgentsRoutes(req: IncomingMessage, res: ServerR const path = url.pathname.slice(CLOUD_AGENTS_PATH.length); try { let data: unknown; - if ((path === '/bootstrap' || path === '/agents') && req.method === 'GET') { + const stream = /^\/runs\/([a-zA-Z0-9_-]{1,64})\/events$/.exec(path); + if (stream && req.method === 'GET') { + const controller = new AbortController(); + const decoder = new TextDecoder(); + const close = () => controller.abort(); + res.once('close', close); + try { + const after = (typeof req.headers['last-event-id'] === 'string' ? req.headers['last-event-id'] : null) + ?? url.searchParams.get('after_seq') ?? '0-0'; + for await (const chunk of cloudAgents.events(stream[1], after, controller.signal)) { + if (!res.headersSent) { + res.setHeader('Content-Type', 'text/event-stream; charset=utf-8'); + res.setHeader('Connection', 'keep-alive'); + flushStreamingHeaders(res); + } + if (!await writeStreamingChunk(res, decoder.decode(chunk, { stream: true }))) break; + } + } finally { + controller.abort(); + res.off('close', close); + } + res.end(); + return true; + } + if (path === '/attachments/pick' && req.method === 'POST') { + data = await cloudAgents.upload(); + } else if (path === '/knowledge/pick' && req.method === 'POST') { + data = await cloudAgents.uploadKnowledge(await parseJsonBody(req)); + } else if (path === '/files/save' && req.method === 'POST') { + data = await cloudAgents.download(await parseJsonBody(req)); + } else if (path === '/actions' && req.method === 'POST') { + data = await cloudAgents.execute(await parseJsonBody(req)); + } else if ((path === '/bootstrap' || path === '/agents') && req.method === 'GET') { data = await cloudAgents.list(url.searchParams.get('cursor')); } else if (path === '/agents' && req.method === 'POST') { data = await cloudAgents.create(await parseJsonBody(req)); @@ -28,7 +60,8 @@ export async function handleCloudAgentsRoutes(req: IncomingMessage, res: ServerR } catch (error) { const failure = error instanceof CloudAgentsError ? error : new CloudAgentsError(error instanceof SyntaxError ? 400 : 502, error instanceof SyntaxError ? 'invalid_input' : 'cloud_service_unavailable'); - sendJson(res, failure.status, { success: false, error: failure.message, code: failure.code }); + if (res.headersSent) res.end(); + else sendJson(res, failure.status, { success: false, error: failure.message, code: failure.code }); } return true; } diff --git a/electron/main/app-deep-link.ts b/electron/main/app-deep-link.ts index bfefe53..90029ea 100644 --- a/electron/main/app-deep-link.ts +++ b/electron/main/app-deep-link.ts @@ -4,7 +4,17 @@ export type NianCodeDeepLink = { type: 'desktop-auth-callback'; requestId: string; url: string; -}; +} | { type: 'cloud-agent'; slug: string; url: string }; + +let pendingCloudAgentRoute: string | null = null; +export function queueCloudAgentLink(slug: string): void { + pendingCloudAgentRoute = '/cloud-agents/shared/' + slug; +} +export function takeCloudAgentRoute(): string | null { + const route = pendingCloudAgentRoute; + pendingCloudAgentRoute = null; + return route; +} const SENSITIVE_QUERY_KEYS = new Set([ 'access_token', @@ -21,6 +31,10 @@ export function parseNianCodeDeepLinkUrl(rawUrl: string): NianCodeDeepLink | nul return null; } + if (url.protocol === `${NIANCODE_APP_PROTOCOL}:` && url.hostname === 'agents' + && /^\/ml-[a-f0-9]{32}$/.test(url.pathname) && !url.username && !url.password && !url.port && !url.search && !url.hash) { + return { type: 'cloud-agent', slug: url.pathname.slice(1), url: rawUrl }; + } if (url.protocol !== `${NIANCODE_APP_PROTOCOL}:`) { return null; } diff --git a/electron/main/index.ts b/electron/main/index.ts index 9dd753e..b958ac7 100644 --- a/electron/main/index.ts +++ b/electron/main/index.ts @@ -29,6 +29,7 @@ import { import { findNianCodeDeepLinkUrl, NIANCODE_APP_PROTOCOL, + queueCloudAgentLink, parseNianCodeDeepLinkUrl, } from './app-deep-link'; import { installMediaPermissionHandler } from './media-permissions'; @@ -355,7 +356,12 @@ function handleAppDeepLinkActivation(rawUrl: string): boolean { return false; } - logger.info(`Received Makelore app link: type=${deepLink.type}, request_id=${deepLink.requestId}`); + if (deepLink.type === 'cloud-agent') { + queueCloudAgentLink(deepLink.slug); + mainWindow?.webContents.send('cloud-agent-link'); + } else { + logger.info(`Received Makelore app link: type=${deepLink.type}, request_id=${deepLink.requestId}`); + } requestMainWindowFocus('app deep link'); return true; } @@ -365,10 +371,7 @@ function registerMakeloreProtocolClient(): void { const registered = devEntrypoint ? app.setAsDefaultProtocolClient(NIANCODE_APP_PROTOCOL, process.execPath, [devEntrypoint]) : app.setAsDefaultProtocolClient(NIANCODE_APP_PROTOCOL); - - if (!registered) { - logger.warn(`Failed to register ${NIANCODE_APP_PROTOCOL}:// protocol handler`); - } + if (!registered) logger.warn(`Failed to register ${NIANCODE_APP_PROTOCOL}:// protocol handler`); } function createMainWindow(): BrowserWindow { diff --git a/electron/main/ipc-handlers.ts b/electron/main/ipc-handlers.ts index e6504cb..78866ba 100644 --- a/electron/main/ipc-handlers.ts +++ b/electron/main/ipc-handlers.ts @@ -21,6 +21,7 @@ import { import type { BackgroundLifecycleController, DesktopActivity } from './background-lifecycle'; import { collectPerformanceSnapshot } from './performance-diagnostics'; import type { HostApiContext } from '../api/context'; +import { takeCloudAgentRoute } from './app-deep-link'; type UnifiedRequest = { id?: string; @@ -261,6 +262,7 @@ export function registerIpcHandlers( ipcMain.handle('app:name', () => app.getName()); ipcMain.handle('app:platform', () => process.platform); ipcMain.handle('app:getPath', (_event, name: Parameters[0]) => app.getPath(name)); + ipcMain.handle('app:take-cloud-agent-link', () => takeCloudAgentRoute()); ipcMain.handle('app:quit', () => app.quit()); ipcMain.handle('app:relaunch', () => { app.relaunch(); diff --git a/electron/preload/index.ts b/electron/preload/index.ts index 62e401e..a500010 100644 --- a/electron/preload/index.ts +++ b/electron/preload/index.ts @@ -23,6 +23,7 @@ const validInvokeChannels = [ 'app:version', 'app:name', 'app:getPath', + 'app:take-cloud-agent-link', 'app:platform', 'app:quit', 'app:relaunch', @@ -64,6 +65,7 @@ const validInvokeChannels = [ const validEventChannels = [ 'navigate', + 'cloud-agent-link', 'update:status-changed', 'update:checking', 'update:available', diff --git a/electron/services/cloud-agent-operations.ts b/electron/services/cloud-agent-operations.ts new file mode 100644 index 0000000..365ac9d --- /dev/null +++ b/electron/services/cloud-agent-operations.ts @@ -0,0 +1,86 @@ +import type { CloudAgentOperations } from '../../shared/cloud-agents'; + +type OperationSpec = { + method: string; path: string; body?: readonly string[]; output: readonly string[]; + query?: Record; cloud?: 'ws'; +}; +const prompt = ['request_id', 'thread_id', 'query', 'expected_revision', 'attachment_file_ids']; +const request = ['request_id', 'thread_id', 'run_id', 'status', 'version']; +const schedule = ['name', 'prompt', 'cron_expression', 'timezone', 'enabled']; +const job = ['id', 'agent_slug', ...schedule, 'next_run_at', 'runs']; +const entry = ['slug', 'name', 'purpose', 'published_version', 'is_creator', 'payer']; +/** Only these product operations can cross the Main boundary; no arbitrary upstream URL or credentials. */ +const operations: Record = { + knowledge: { method: 'GET', path: '/agents/:slug/knowledge', output: ['databases', 'models'] }, + createKnowledge: { method: 'POST', path: '/agents/:slug/knowledge', body: ['operation_id', 'name', 'embedding_model'], output: ['kb_id', 'name', 'description', 'embedding_model'] }, + knowledgeFiles: { method: 'GET', path: '/agents/:slug/knowledge/:kb_id/files', query: { offset: 'offset' }, output: ['files', 'next_offset'] }, + processKnowledge: { method: 'POST', path: '/agents/:slug/knowledge/:kb_id/files/:file_id/process', body: ['operation_id'], output: ['task_id'] }, + catalog: { method: 'GET', path: '/catalog', output: ['models', 'resources', 'pricing'] }, + received: { method: 'GET', path: '/received', output: ['agents'] }, + entry: { method: 'GET', path: '/entry/:slug', output: entry }, + publish: { method: 'POST', path: '/agents/:slug/publish', body: ['operation_id', 'expected_revision'], output: ['version', 'draft_revision'] }, + access: { method: 'GET', path: '/agents/:slug/access', output: ['published_version', 'enabled', 'share_url', 'versions', 'grants', 'applications'] }, + setEnabled: { method: 'PATCH', path: '/agents/:slug/enabled', body: ['enabled'], output: ['enabled'] }, + users: { method: 'GET', cloud: 'ws', path: '/users', query: { query: 'query' }, output: ['users'] }, + share: { method: 'PUT', path: '/agents/:slug/shares/:account_id', body: ['enabled'], output: ['account_id', 'enabled'] }, + createApplication: { method: 'POST', path: '/agents/:slug/applications', body: ['operation_id', 'name'], output: ['id', 'name', 'enabled'] }, + setApplicationEnabled: { method: 'PATCH', path: '/applications/:application_id', body: ['enabled'], output: ['enabled'] }, + keys: { method: 'GET', path: '/applications/:application_id/keys', output: ['keys'] }, + createKey: { method: 'POST', path: '/applications/:application_id/keys', body: ['operation_id'], output: ['id', 'prefix', 'secret'] }, + revokeKey: { method: 'DELETE', path: '/applications/:application_id/keys/:key_id', output: ['revoked'] }, + costs: { method: 'GET', cloud: 'ws', path: '/costs', query: { slug: 'agent_slug' }, output: ['unit', 'items'] }, + threads: { method: 'GET', path: '/threads', query: { slug: 'slug', offset: 'offset' }, output: ['threads', 'next_offset'] }, + createThread: { method: 'POST', path: '/agents/:slug/threads', body: ['thread_id', 'preview', 'expected_revision'], output: ['thread_id', 'client_thread_id'] }, + history: { method: 'GET', path: '/threads/:thread_id', query: { offset: 'offset' }, output: ['thread_id', 'messages', 'run', 'next_offset'] }, + viewed: { method: 'POST', path: '/threads/:thread_id/viewed', body: ['run_id'], output: ['viewed'] }, + submit: { method: 'POST', path: '/agents/:slug/requests', body: prompt, output: request }, + preview: { method: 'POST', path: '/agents/:slug/preview', body: prompt, output: request }, + request: { method: 'GET', path: '/requests/:request_id', output: request }, + cancelRequest: { method: 'POST', path: '/requests/:request_id/cancel', output: ['status'] }, + run: { method: 'GET', path: '/runs/:run_id', output: ['agent_run_id', 'request_id', 'thread_id', 'agent_slug', 'status', 'output', 'version', 'error', 'interrupt'] }, + cancelRun: { method: 'POST', path: '/runs/:run_id/cancel', output: ['status'] }, + resume: { method: 'POST', path: '/runs/:run_id/resume', body: ['operation_id', 'decision'], output: ['run_id', 'status'] }, + schedules: { method: 'GET', path: '/agents/:slug/schedules', output: ['jobs'] }, + createSchedule: { method: 'POST', path: '/agents/:slug/schedules', body: ['operation_id', ...schedule], output: job }, + updateSchedule: { method: 'PUT', path: '/agents/:slug/schedules/:job_id', body: schedule, output: job }, + deleteSchedule: { method: 'DELETE', path: '/agents/:slug/schedules/:job_id', output: ['deleted'] }, + runSchedule: { method: 'POST', path: '/agents/:slug/schedules/:job_id/run-now', body: ['operation_id'], output: ['thread_id', 'status'] }, + attachments: { method: 'GET', path: '/threads/:thread_id/attachments', output: ['attachments'] }, + parseAttachment: { method: 'POST', path: '/attachments/tmp/parse', body: ['object_name', 'parse_method'], output: ['parsed_object_name'] }, + confirmAttachment: { method: 'POST', path: '/threads/:thread_id/attachments/confirm', body: ['attachments'], output: ['attachments'] }, + deleteAttachment: { method: 'DELETE', path: '/threads/:thread_id/attachments/:file_id', output: ['message'] }, + files: { method: 'GET', path: '/threads/:thread_id/files', query: { path: 'path' }, output: ['files'] }, +}; + +export function operationPlan(value: unknown) { + if (!value || typeof value !== 'object') throw new Error('invalid_operation'); + const { operation, input } = value as { operation?: unknown; input?: unknown }; + if (typeof operation !== 'string' || !Object.hasOwn(operations, operation) + || !input || typeof input !== 'object' || Array.isArray(input)) throw new Error('invalid_operation'); + const spec = operations[operation as keyof CloudAgentOperations]; + const args = input as Record; + const path = spec.path.replace(/:([a-z_]+)/g, (_, key: string) => { + const id = args[key]; + if (typeof id !== 'string' || !/^[a-zA-Z0-9_-]{1,128}$/.test(id) + || (key === 'slug' && !/^ml-[a-f0-9]{32}$/.test(id))) throw new Error('invalid_id'); + return encodeURIComponent(id); + }); + const query = new URLSearchParams(); + for (const [key, name] of Object.entries(spec.query ?? {})) { + const v = args[key]; + if (v === undefined) continue; + if (key === 'offset' ? !Number.isSafeInteger(v) || Number(v) < 0 : typeof v !== 'string' || v.length > (key === 'path' ? 2048 : 100)) { + throw new Error('invalid_query'); + } + query.set(name, String(v)); + } + return { + ...spec, path: (spec.cloud === 'ws' ? '/api/cloud-agents' : '/api/makelore') + path + (query.size ? '?' + query : ''), + body: spec.body ? Object.fromEntries(spec.body.filter(key => args[key] !== undefined).map(key => [key, args[key]])) : undefined, + project: (result: unknown) => { + if (!result || typeof result !== 'object' || Array.isArray(result)) throw new Error('invalid_response'); + const record = result as Record; + return Object.fromEntries(spec.output.filter(key => record[key] !== undefined).map(key => [key, record[key]])); + }, + }; +} diff --git a/electron/services/cloud-agents.ts b/electron/services/cloud-agents.ts index 97cba47..99b53f7 100644 --- a/electron/services/cloud-agents.ts +++ b/electron/services/cloud-agents.ts @@ -1,4 +1,5 @@ -import type { CloudAgentDraft, CloudAgentPage, CreateCloudAgent, SaveCloudAgentDraft } from '../../shared/cloud-agents'; +import { EMPTY_CLOUD_CONFIGURATION, type CloudAgentConfiguration, type CloudAgentDraft, type CloudAgentPage, type CreateCloudAgent, type SaveCloudAgentDraft, type CloudUpload } from '../../shared/cloud-agents'; +import { operationPlan } from './cloud-agent-operations'; import { WORKS_SQUARE_CONFIG } from '../api/works-config'; import { proxyAwareFetch, runWithDeadline } from '../utils/proxy-fetch'; import { @@ -19,6 +20,17 @@ const MESSAGES: Record = { operation_conflict: '创建操作的内容发生冲突,请刷新列表确认', cloud_service_unavailable: '智能体服务暂时不可用,请重试', invalid_input: '请检查名称、用途和配置内容', + model_required: '请先选择模型并保存', + model_unavailable: '所选模型当前不可用,请刷新模型目录', + resource_unavailable: '所选资源已不可用,请刷新配置', + publish_required: '请先发布智能体', + access_revoked: '该智能体的使用授权已撤销', + application_disabled: '应用已停用', + key_revoked: '此凭据已撤销,请新建凭据', + insufficient_balance: '创建者的词元点数不足', + request_conflict: '该请求已用于其他输入,请开始新的请求', + attachment_too_large: '附件最大支持 5 MB', + download_failed: '文件保存失败,请重试', }; export class CloudAgentsError extends Error { @@ -60,9 +72,32 @@ function projectDraft(value: unknown): CloudAgentDraft { system_prompt: textField(item.system_prompt, 32000, false), draft_revision: revision(item.draft_revision), updated_at: textField(item.updated_at, 40), + configuration: configuration(item.configuration), + published_version: item.published_version == null ? null : revision(item.published_version), + enabled: item.enabled !== false, }; } catch { throw new CloudAgentsError(502, 'cloud_service_unavailable'); } } +function configuration(value: unknown): CloudAgentConfiguration { + const input = record(value); + const result = { ...EMPTY_CLOUD_CONFIGURATION, ...Object.fromEntries( + Object.keys(EMPTY_CLOUD_CONFIGURATION).filter(key => input[key] !== undefined).map(key => [key, input[key]]), + ) }; + textField(result.model, 200, false); + for (const key of ['tools', 'knowledges', 'mcps', 'skills', 'preload_skills', 'subagents'] as const) { + if (!Array.isArray(result[key]) || result[key].length > 100 || result[key].some(v => typeof v !== 'string' || v.length > 200)) { + throw new CloudAgentsError(422, 'invalid_input'); + } + } + if (!['default', 'always_trust'].includes(result.tool_approval_mode) + || !Number.isInteger(result.max_execution_steps) || result.max_execution_steps < 1 || result.max_execution_steps > 300 + || !Number.isInteger(result.max_output_tokens) || result.max_output_tokens < 1 || result.max_output_tokens > 32768 + || !Number.isInteger(result.max_run_seconds) || result.max_run_seconds < 10 || result.max_run_seconds > 3600) { + throw new CloudAgentsError(422, 'invalid_input'); + } + return result; +} + type Session = { binding: WorksSquareAccountBinding; accessToken: string; @@ -70,19 +105,33 @@ type Session = { expiresAt: number; }; -/** Main owns both clouds' credentials; public methods return only Agent drafts. */ +/** Main owns cloud credentials and transport; Renderer receives product data only. */ export class CloudAgentsModule { private cached: Session | null = null; private readonly unsubscribe: () => void; + private readonly streams = new Map(); + private activityTimer: ReturnType; + private activityPolling = false; + private activitySeen: Map | null = null; constructor(private readonly fetchImpl: (input: string, init?: RequestInit) => Promise = proxyAwareFetch) { this.unsubscribe = subscribeWorksSquareSession(() => { - if (this.cached && !isCurrentWorksSquareAccountBinding(this.cached.binding)) this.cached = null; + if (this.cached && !isCurrentWorksSquareAccountBinding(this.cached.binding)) { + this.cached = null; + this.activitySeen = null; + } + for (const [controller, binding] of this.streams) { + if (!isCurrentWorksSquareAccountBinding(binding)) controller.abort(); + } }); + this.activityTimer = setInterval(() => { void this.pollActivity(); }, 30000); + this.activityTimer.unref?.(); } dispose(): void { this.cached = null; + clearInterval(this.activityTimer); + for (const controller of this.streams.keys()) controller.abort(); this.unsubscribe(); } @@ -121,21 +170,210 @@ export class CloudAgentsModule { name: textField(input.name, 100).trim(), purpose: textField(input.purpose, 2000).trim(), system_prompt: textField(input.system_prompt, 32000, false), + ...(input.configuration === undefined ? {} : { configuration: configuration(input.configuration) }), }; return projectDraft(await this.request(`/api/makelore/agents/${slug(agentSlug)}/draft`, 'PATCH', body)); } - private async request(path: string, method: string, body?: unknown): Promise { + private async pollActivity(): Promise { + if (!this.cached || this.activityPolling) return; + this.activityPolling = true; + try { + const binding = this.cached.binding; + const result = record(await this.request('/api/makelore/threads', 'GET')); + const rows = Array.isArray(result.threads) ? result.threads.map(record) : []; + const next = new Map(rows.filter(row => typeof row.run_id === 'string').map(row => [String(row.run_id), String(row.status)])); + if (this.activitySeen) { + const changed = rows.filter(row => row.unread === true && typeof row.run_id === 'string' + && ['completed', 'failed', 'interrupted'].includes(String(row.status)) + && this.activitySeen?.get(row.run_id) !== row.status); + if (changed.length) { + const { Notification, BrowserWindow } = await import('electron'); + this.requireCurrent(binding); + if (Notification.isSupported()) { + const notice = new Notification({ title: '智能体有新的任务结果', body: changed.some(row => row.status === 'interrupted') + ? '有任务等待你的确认,打开活动查看。' : '打开活动查看完成结果或失败原因。' }); + notice.on('click', () => { + if (!isCurrentWorksSquareAccountBinding(binding)) return; + const window = BrowserWindow.getAllWindows().find(w => !w.isDestroyed()); + window?.show(); window?.focus(); window?.webContents.send('navigate', '/cloud-agents?view=activity'); + }); + notice.show(); + } + } + } + this.activitySeen = next; + } catch(error) { + if (error instanceof CloudAgentsError && [401, 403, 423].includes(error.status)) this.cached = null; + // The durable activity list remains the source of truth when notifications or transport are unavailable. + } finally { this.activityPolling = false; } + } + + async execute(value: unknown): Promise { + let plan; + try { plan = operationPlan(value); } + catch { throw new CloudAgentsError(422, 'invalid_input'); } + const result = await this.request(plan.path, plan.method, plan.body, plan.cloud === 'ws'); + try { + const projected = plan.project(result); + if (record(value).operation === 'access') { + // The address is public; both exchanged tokens stay in Main. + const binding = getWorksSquareAccountBinding(); + if (!binding) throw new CloudAgentsError(401, 'session_expired'); + const session = await runWithDeadline(signal => this.session(binding, signal), 30000); + this.requireCurrent(binding); + return { ...projected, api_url: session.apiBaseUrl + '/api/makelore/api/requests' }; + } + return projected; + } + catch (error) { + if (error instanceof CloudAgentsError) throw error; + throw new CloudAgentsError(502, 'cloud_service_unavailable'); + } + } + + async upload(): Promise { + const body = await this.pickUpload(); + if (!body) return null; + const result = record(await this.request('/api/makelore/attachments/tmp', 'POST', body)); + return { + object_name: textField(result.object_name, 2048), file_name: textField(result.file_name, 255), + file_type: typeof result.file_type === 'string' ? result.file_type : 'application/octet-stream', + parse_supported: result.parse_supported === true, + parse_methods: Array.isArray(result.parse_methods) ? result.parse_methods.filter((v): v is string => typeof v === 'string') : [], + }; + } + + async uploadKnowledge(value: unknown): Promise { + const input = record(value); + const slug = textField(input.slug, 128), kbId = textField(input.kb_id, 128); + const operationId = textField(input.operation_id, 64); + if (!/^ml-[a-f0-9]{32}$/.test(slug) || !/^kb_[a-z0-9]+$/.test(kbId) + || !/^[a-f0-9-]{36}$/.test(operationId)) throw new CloudAgentsError(422, 'invalid_input'); + const body = await this.pickUpload(); + if (!body) return null; + const result = record(await this.request(`/api/makelore/agents/${slug}/knowledge/${kbId}/files?operation_id=${operationId}`, 'POST', body)); + return Object.fromEntries(['file_id', 'name', 'size', 'status', 'error', 'chunk_count'].map(key => [key, result[key]])); + } + + private async pickUpload(): Promise { + const binding = getWorksSquareAccountBinding(); + if (!binding) throw new CloudAgentsError(401, 'session_expired'); + const { dialog } = await import('electron'); + const { readFile, stat } = await import('node:fs/promises'); + const { basename } = await import('node:path'); + const picked = await dialog.showOpenDialog({ title: '添加智能体附件(最大 5 MB)', properties: ['openFile'] }); + this.requireCurrent(binding); + if (picked.canceled || !picked.filePaths[0]) return null; + const path = picked.filePaths[0]; + if ((await stat(path)).size > 5 * 1024 * 1024) throw new CloudAgentsError(422, 'attachment_too_large'); + const bytes = await readFile(path); + this.requireCurrent(binding); + if (bytes.length > 5 * 1024 * 1024) throw new CloudAgentsError(422, 'attachment_too_large'); + const body = new FormData(); + body.append('file', new Blob([new Uint8Array(bytes)]), basename(path)); + return body; + } + + async download(value: unknown): Promise<{ saved: boolean }> { + const input = record(value); + const threadId = textField(input.thread_id, 128); + const path = textField(input.path, 2048); + if (!/^[a-zA-Z0-9_-]+$/.test(threadId)) throw new CloudAgentsError(422, 'invalid_input'); + const binding = getWorksSquareAccountBinding(); + if (!binding) throw new CloudAgentsError(401, 'session_expired'); + const { dialog } = await import('electron'); + const fs = await import('node:fs/promises'); + const { basename, dirname, join } = await import('node:path'); + const picked = await dialog.showSaveDialog({ title: '保存智能体文件', defaultPath: basename(path), + properties: ['createDirectory', 'showOverwriteConfirmation'] }); + this.requireCurrent(binding); + if (picked.canceled || !picked.filePath) return { saved: false }; + const temporary = join(dirname(picked.filePath), '.makelore-download-' + crypto.randomUUID()); + try { + await runWithDeadline(async signal => { + const session = await this.session(binding, signal); + this.requireCurrent(binding); + const url = session.apiBaseUrl + '/api/makelore/threads/' + encodeURIComponent(threadId) + '/artifacts/' + + path.replace(/^\/+/, '').split('/').map(encodeURIComponent).join('/') + '?download=true'; + const response = await this.fetchImpl(url, { redirect: 'error', signal, headers: { Authorization: 'Bearer ' + session.accessToken } }); + if (!response.ok || !response.body) throw new CloudAgentsError(response.status || 502, 'download_failed'); + const output = await fs.open(temporary, 'wx'); + const reader = response.body.getReader(); + let size = 0; + try { + while (true) { + const { value: chunk, done } = await reader.read(); + this.requireCurrent(binding); + if (done) break; + size += chunk.byteLength; + if (size > 1024 * 1024 * 1024) throw new CloudAgentsError(413, 'download_failed'); + let offset = 0; + while (offset < chunk.byteLength) offset += (await output.write(chunk, offset, chunk.byteLength - offset)).bytesWritten; + } + } finally { await reader.cancel().catch(() => undefined); await output.close(); } + }, 300000); + this.requireCurrent(binding); + await fs.rename(temporary, picked.filePath); + return { saved: true }; + } catch(error) { + if (error instanceof CloudAgentsError) throw error; + throw new CloudAgentsError(502, 'download_failed'); + } finally { await fs.unlink(temporary).catch(() => undefined); } + } + + async *events(runId: string, after: string, signal: AbortSignal): AsyncGenerator { + if (!/^[a-zA-Z0-9_-]{1,64}$/.test(runId) || !/^\d+-\d+$/.test(after)) { + throw new CloudAgentsError(422, 'invalid_input'); + } + const binding = getWorksSquareAccountBinding(); + if (!binding) throw new CloudAgentsError(401, 'session_expired'); + const controller = new AbortController(); + const abort = () => controller.abort(); + signal.addEventListener('abort', abort, { once: true }); + if (signal.aborted) abort(); + this.streams.set(controller, binding); + let cancelReader: (() => Promise) | undefined; + try { + const session = await runWithDeadline(s => this.session(binding, s), 30000, controller.signal); + this.requireCurrent(binding); + const response = await this.fetchImpl(session.apiBaseUrl + '/api/makelore/runs/' + runId + '/events?after_seq=' + after, { + method: 'GET', redirect: 'error', signal: controller.signal, + headers: { Authorization: 'Bearer ' + session.accessToken, Accept: 'text/event-stream' }, + }); + this.requireCurrent(binding); + if (!response.ok || !response.body) throw new CloudAgentsError(response.status || 502, 'cloud_service_unavailable'); + const reader = response.body.getReader(); + cancelReader = () => reader.cancel(); + while (true) { + const chunk = await reader.read(); + this.requireCurrent(binding); + if (chunk.done) break; + yield chunk.value; + } + } finally { + await cancelReader?.().catch(() => undefined); + controller.abort(); + signal.removeEventListener('abort', abort); + this.streams.delete(controller); + } + } + + private async request(path: string, method: string, body?: unknown, worksSquare = false): Promise { const binding = getWorksSquareAccountBinding(); if (!binding) throw new CloudAgentsError(401, 'session_expired'); try { return await runWithDeadline(async (signal) => { - const session = await this.session(binding, signal); + const session = worksSquare ? { + apiBaseUrl: WORKS_SQUARE_CONFIG.apiBaseUrl, + accessToken: await getValidWorksSquareAccessToken(), + } : await this.session(binding, signal); this.requireCurrent(binding); + if (!session.accessToken) throw new CloudAgentsError(401, 'session_expired'); const result = await this.fetchJson(session.apiBaseUrl + path, { method, signal, redirect: 'error', - headers: { Authorization: `Bearer ${session.accessToken}`, 'Content-Type': 'application/json' }, - ...(body === undefined ? {} : { body: JSON.stringify(body) }), + headers: { Authorization: `Bearer ${session.accessToken}`, ...(body instanceof FormData ? {} : { 'Content-Type': 'application/json' }) }, + ...(body === undefined ? {} : { body: body instanceof FormData ? body : JSON.stringify(body) }), }); this.requireCurrent(binding); return result; @@ -144,7 +382,6 @@ export class CloudAgentsModule { this.requireCurrent(binding); if (error instanceof CloudAgentsError) { if (error.status === 401) this.cached = null; - if (error.status === 422) throw new CloudAgentsError(502, 'cloud_service_unavailable'); throw error; } throw new CloudAgentsError(502, 'cloud_service_unavailable'); @@ -171,7 +408,7 @@ export class CloudAgentsModule { const url = new URL(textField(value.api_base_url, 2048)); if ((url.protocol !== 'https:' && !(url.protocol === 'http:' && ['127.0.0.1', 'localhost', '[::1]'].includes(url.hostname))) || url.username || url.password || url.search || url.hash - || value.scope !== 'makelore-agent-drafts' || value.token_type !== 'bearer' + || value.scope !== 'makelore-agents' || value.token_type !== 'bearer' || typeof value.expires_at !== 'number' || value.expires_at * 1000 <= Date.now() || value.expires_at * 1000 > Date.now() + 300000) { throw new CloudAgentsError(502, 'cloud_service_unavailable'); @@ -191,7 +428,9 @@ export class CloudAgentsModule { if (!response.ok) { const code = record(record(body).detail).code; throw new CloudAgentsError( - response.status, typeof code === 'string' && Object.hasOwn(MESSAGES, code) ? code : 'cloud_service_unavailable', + response.status, typeof code === 'string' && Object.hasOwn(MESSAGES, code) ? code + : response.status === 422 ? 'invalid_input' : response.status === 401 ? 'session_expired' + : response.status === 403 ? 'access_revoked' : 'cloud_service_unavailable', ); } return body; diff --git a/shared/cloud-agents.ts b/shared/cloud-agents.ts index 7c58669..1286730 100644 --- a/shared/cloud-agents.ts +++ b/shared/cloud-agents.ts @@ -7,6 +7,9 @@ export interface CloudAgentDraft { system_prompt: string; draft_revision: number; updated_at: string; + configuration: CloudAgentConfiguration; + published_version: number | null; + enabled: boolean; } export interface CloudAgentPage { @@ -25,4 +28,137 @@ export interface SaveCloudAgentDraft { name: string; purpose: string; system_prompt: string; + configuration?: CloudAgentConfiguration; +} + +export interface CloudAgentConfiguration { + model: string; + tools: string[]; + knowledges: string[]; + mcps: string[]; + skills: string[]; + preload_skills: string[]; + subagents: string[]; + tool_approval_mode: 'default' | 'always_trust'; + max_execution_steps: number; + max_output_tokens: number; + max_run_seconds: number; +} + +export const EMPTY_CLOUD_CONFIGURATION: CloudAgentConfiguration = { + model: '', tools: [], knowledges: [], mcps: [], skills: [], preload_skills: [], subagents: [], + tool_approval_mode: 'default', max_execution_steps: 40, max_output_tokens: 4096, max_run_seconds: 600, +}; + +export interface CloudAgentEntry { + slug: string; name: string; purpose: string; published_version: number; + is_creator?: boolean; payer?: string; +} +export interface CloudResource { key: string; name: string; description?: string } +export interface CloudCatalog { + models: { id: string; name: string }[]; + resources: Partial>; + pricing: { unit: string; usage_unit: string; unit_size: number; points_per_unit: string; version: number }; +} +export interface CloudApplication { id: string; name: string; enabled: boolean } +export interface CloudKey { id: string; prefix: string; revoked_at: string | null } +export interface CloudAccess { + published_version: number | null; enabled: boolean; share_url: string; + versions: { version: number; draft_revision: number; created_at: string }[]; + grants: { account_id: string; enabled: boolean }[]; + applications: CloudApplication[]; + api_url: string; +} +export interface CloudThread { + thread_id: string; client_thread_id: string | null; agent_slug: string; title: string; + mode: 'published' | 'preview'; updated_at: string; run_id: string | null; status: string; unread: boolean; + draft_revision?: number; +} +export interface CloudRun { + agent_run_id: string; request_id: string; thread_id: string; agent_slug: string; + status: string; output: string; version: string; + error?: { type: string; message: string }; + interrupt?: CloudInterrupt; +} +export interface CloudInterrupt { + status?: string; message?: string; run_id?: string; + questions?: unknown[]; approval?: Record; interrupt_info?: unknown; +} +export interface CloudMessage { + id: number; role: 'user' | 'assistant' | 'system'; content: string; request_id: string; + run_id: string | null; created_at: string; +} +export interface CloudHistory { + thread_id: string; messages: CloudMessage[]; run: CloudRun | null; next_offset: number | null; +} +export interface CloudRequest { + request_id: string; thread_id: string; run_id: string | null; status: string; version: string; +} +export interface CloudPrompt { + request_id: string; thread_id: string; query: string; expected_revision?: number; + attachment_file_ids?: string[]; +} +export interface CloudScheduleInput { + name: string; prompt: string; cron_expression: string; timezone: string; enabled: boolean; +} +export interface CloudSchedule extends CloudScheduleInput { + id: string; agent_slug: string; next_run_at: string | null; + runs?: { status: string; thread_id: string; error_message?: string; conversation_available: boolean }[]; +} +export interface CloudCost { + id: string; status: string; reserved_points: string; actual_points: string | null; created_at: string; + context: { agent_slug: string; version: string; caller_kind: string; caller_id: string; run_id: string; request_id: string }; +} +export interface CloudAttachment { file_id: string; file_name: string; file_size: number; path: string; original_path: string; request_id?: string | null } +export interface CloudUpload { object_name: string; file_name: string; file_type: string; parse_supported: boolean; parse_methods: string[] } +export interface CloudFile { name: string; path: string; directory_path: string; is_dir: boolean; size: number } +export interface CloudKnowledge { kb_id: string; name: string; description: string; embedding_model: string } +export interface CloudKnowledgeFile { file_id: string; name: string; size: number; status: string; error: string | null; chunk_count: number } +type Op = { input: Input; output: Output }; +type Agent = { slug: string }; +type Application = { application_id: string }; +type Schedule = Agent & { job_id: string }; +type Thread = { thread_id: string }; +type Run = { run_id: string }; +type Operation = { operation_id: string }; +export interface CloudAgentOperations { + knowledge: Op; + createKnowledge: Op; + knowledgeFiles: Op; + processKnowledge: Op; + catalog: Op, CloudCatalog>; + received: Op, { agents: CloudAgentEntry[] }>; + entry: Op; + publish: Op; + access: Op; + setEnabled: Op; + users: Op<{ query: string }, { users: { account_id: string; username: string; display_name: string }[] }>; + share: Op; + createApplication: Op; + setApplicationEnabled: Op; + keys: Op; + createKey: Op; + revokeKey: Op; + costs: Op<{ slug?: string }, { unit: string; items: CloudCost[] }>; + threads: Op<{ slug?: string; offset?: number }, { threads: CloudThread[]; next_offset: number | null }>; + createThread: Op; + history: Op; + viewed: Op; + submit: Op; + preview: Op; + request: Op<{ request_id: string }, CloudRequest>; + cancelRequest: Op<{ request_id: string }, { status: string }>; + run: Op; + cancelRun: Op; + resume: Op }, { run_id: string; status: string }>; + schedules: Op; + createSchedule: Op; + updateSchedule: Op; + deleteSchedule: Op; + runSchedule: Op; + attachments: Op; + parseAttachment: Op<{ object_name: string; parse_method: string }, { parsed_object_name: string }>; + confirmAttachment: Op; + deleteAttachment: Op; + files: Op; } diff --git a/src/App.tsx b/src/App.tsx index 55970cf..6395e60 100644 --- a/src/App.tsx +++ b/src/App.tsx @@ -24,6 +24,7 @@ import { useUserSyncStore } from './stores/user-sync'; import { flushPendingAgentSessionSync } from '@/lib/agent-session-sync'; import { subscribeHostEvent } from '@/lib/host-events'; import { reportDesktopActivity, type DesktopActivityModule } from '@/lib/host-api'; +import { invokeIpc } from '@/lib/api-client'; import { installRendererPerformanceDiagnostics } from '@/lib/performance-diagnostics'; import { resolveSupportedLanguage } from '../shared/language'; import { @@ -399,11 +400,19 @@ function App() { }; const unsubscribe = window.electron.ipcRenderer.on('navigate', handleNavigate); + const openAgentLink = () => { + void invokeIpc('app:take-cloud-agent-link').then(path => { + if (path && /^\/cloud-agents\/shared\/ml-[a-f0-9]{32}$/.test(path)) navigate(path); + }).catch(() => undefined); + }; + const unsubscribeAgentLink = window.electron.ipcRenderer.on('cloud-agent-link', openAgentLink); + openAgentLink(); return () => { if (typeof unsubscribe === 'function') { unsubscribe(); } + if (typeof unsubscribeAgentLink === 'function') unsubscribeAgentLink(); }; }, [navigate]); diff --git a/src/lib/cloud-agents-api.ts b/src/lib/cloud-agents-api.ts index 42c00ff..79b85ec 100644 --- a/src/lib/cloud-agents-api.ts +++ b/src/lib/cloud-agents-api.ts @@ -1,7 +1,24 @@ -import { hostApiFetch } from './host-api'; +import { hostApiFetch, createHostEventSource, ensureHostApiToken } from './host-api'; +import type { CloudAgentOperations, CloudUpload, CloudKnowledgeFile } from '../../shared/cloud-agents'; import { CLOUD_AGENTS_PATH, type CloudAgentDraft, type CloudAgentPage, type CreateCloudAgent, type SaveCloudAgentDraft } from '../../shared/cloud-agents'; export const cloudAgentsApi = { + uploadKnowledge: (slug: string, kb_id: string, operation_id: string) => + hostApiFetch(CLOUD_AGENTS_PATH + '/knowledge/pick', { + method: 'POST', body: JSON.stringify({ slug, kb_id, operation_id }), + }), + upload: () => hostApiFetch(CLOUD_AGENTS_PATH + '/attachments/pick', { method: 'POST' }), + download: (thread_id: string, path: string) => hostApiFetch<{ saved: boolean }>(CLOUD_AGENTS_PATH + '/files/save', { + method: 'POST', body: JSON.stringify({ thread_id, path }), + }), + call: (operation: K, input: CloudAgentOperations[K]['input']) => + hostApiFetch(CLOUD_AGENTS_PATH + '/actions', { + method: 'POST', body: JSON.stringify({ operation, input }), + }), + events: async (runId: string, after = '0-0') => { + await ensureHostApiToken(); + return createHostEventSource(CLOUD_AGENTS_PATH + '/runs/' + encodeURIComponent(runId) + '/events?after_seq=' + encodeURIComponent(after)); + }, list: (cursor: string | null = null) => hostApiFetch( CLOUD_AGENTS_PATH + '/agents' + (cursor ? `?cursor=${encodeURIComponent(cursor)}` : ''), ), diff --git a/src/pages/CloudAgents/CloudAccess.tsx b/src/pages/CloudAgents/CloudAccess.tsx new file mode 100644 index 0000000..f2b8d09 --- /dev/null +++ b/src/pages/CloudAgents/CloudAccess.tsx @@ -0,0 +1,144 @@ +import { useCallback, useEffect, useRef, useState } from 'react'; +import { Copy, Link as LinkIcon, Plus } from 'lucide-react'; +import { Button } from '@/components/ui/button'; +import { Input } from '@/components/ui/input'; +import { usePendingCloudInput } from './CloudPending'; +import { cloudAgentsApi } from '@/lib/cloud-agents-api'; +import type { CloudAccess as Access, CloudApplication, CloudCost, CloudKey } from '../../../shared/cloud-agents'; + +const errorText = (e: unknown) => e instanceof Error ? e.message : '操作失败,请重试'; +export function CloudAccessPanel({ slug, revision, onPublished }: { slug: string; revision: number; onPublished: () => void }) { + const [access, setAccess] = useState(null); + const [costs, setCosts] = useState([]); + const [query, setQuery] = useState(''); + const [users, setUsers] = useState<{ account_id: string; display_name: string; username: string }[]>([]); + const [name, setName] = useState(''); + const [busy, setBusy] = useState(false); + const [error, setError] = useState(''); + const [notice, setNotice] = useState(''); + usePendingCloudInput(busy || Boolean(name.trim())); + const publishOperation = useRef<{ operation_id: string; expected_revision: number } | null>(null); + const applicationOperation = useRef<{ operation_id: string; name: string } | null>(null); + const alive = useRef(true); + const refresh = useCallback(async () => { + const results = await Promise.allSettled([cloudAgentsApi.call('access', { slug }), cloudAgentsApi.call('costs', { slug })]); + if (!alive.current) return; + if (results[0].status === 'fulfilled') setAccess(results[0].value); + if (results[1].status === 'fulfilled') setCosts(results[1].value.items); + const failure = results.find(result => result.status === 'rejected'); + if (failure?.status === 'rejected') throw failure.reason; + }, [slug]); + useEffect(() => { + alive.current = true; + refresh().catch(e => { if (alive.current) setError(errorText(e)); }); + return () => { alive.current = false; }; + }, [refresh]); + const act = async (task: () => Promise, success?: string) => { + if (busy) return; + setBusy(true); setError(''); setNotice(''); + try { await task(); await refresh(); if (alive.current && success) setNotice(success); } + catch(e) { if (alive.current) setError(errorText(e)); } + finally { if (alive.current) setBusy(false); } + }; + const publish = () => act(async () => { + const input = publishOperation.current ?? { operation_id: crypto.randomUUID(), expected_revision: revision }; + publishOperation.current = input; + await cloudAgentsApi.call('publish', { slug, ...input }); + publishOperation.current = null; + onPublished(); + }, '已发布,新的调用将使用这个版本'); + const createApplication = () => act(async () => { + const input = applicationOperation.current ?? { operation_id: crypto.randomUUID(), name: name.trim() }; + applicationOperation.current = input; + await cloudAgentsApi.call('createApplication', { slug, ...input }); + applicationOperation.current = null; + setName(''); + }, '应用已创建,可为它生成 API 凭据'); + const copy = async (text: string) => { + try { await navigator.clipboard.writeText(text); setNotice('已复制'); } + catch { setError('复制失败,请手动选择文本复制'); } + }; + return
+
+
+

{access?.published_version ? '当前发布版本 ' + access.published_version : '尚未发布'}

+

发布保存的修订 {revision}。发布后仅自己可用,分享与应用需要分别开启。

+
+

自己使用、分享使用、API 调用和自动任务产生的费用,均从你的个人词元点数扣除。你可以随时停用智能体或撤销访问。

+ {access?.published_version &&
{access.enabled ? '智能体已启用' : '智能体已停用'} +
} +
+ {error &&

{error}

} + {notice &&

{notice}

} +

分享给指定用户

+

选择获准使用的账号,再发送分享链接。每个人的对话和文件分别保存。

+
setQuery(e.target.value)} placeholder="用户名、显示名或账号 ID" maxLength={100} /> +
+ {users.map(user =>
+

{user.display_name} @{user.username}

{user.account_id}

+ +
)} + {access?.grants.filter(g => g.enabled).map(grant =>
+ {grant.account_id} +
)} + {access?.published_version &&
+ {access.share_url} +
} +
+

应用与 API

+

为网站或脚本创建独立应用。应用通过凭据调用此智能体,并拥有自己的会话空间。

+
setName(e.target.value)} maxLength={100} placeholder="例如:我的个人网站" /> +
+ {access?.applications.map(app => )} + {access?.api_url &&
API 调用示例 +

POST {access.api_url}

+
{'Authorization: Bearer <应用凭据>\nContent-Type: application/json\n\n' + JSON.stringify({ request_id: '唯一请求标识', thread_id: '会话标识', query: '帮我整理今天的创作计划' }, null, 2)}
+

相同请求标识与输入可重试。通过返回的 request_id 查询排队状态,通过 run_id 读取运行结果或订阅事件;这些请求继续使用同一应用凭据。

+
} +
+

最近费用

+

最近 100 次模型调用,仅显示费用归属。调用者的对话与文件不会在这里展示。

+ {!costs.length ?

还没有模型调用费用

:
+ + {costs.map(cost => + + )} +
时间调用来源状态实际点数
{new Date(cost.created_at).toLocaleString()}

{cost.context.caller_kind === 'application' ? '应用' : '个人账号'} · v{cost.context.version}

{cost.context.caller_id}

{({ reserved: '已预占', dispatched: '待结算', pending_review: '待核对用量', settled: '已结算', released: '已释放', failed: '失败' } as Record)[cost.status] ?? cost.status}{cost.actual_points === null ? '待结算' : cost.actual_points}
} +
+ {access &&
发布历史({access.versions.length}) + {access.versions.map(v =>

版本 {v.version} · 草稿修订 {v.draft_revision} · {new Date(v.created_at).toLocaleString()}

)}
} +
; +} + +function ApplicationKeys({ application, onUpdated }: { application: CloudApplication; onUpdated: () => Promise }) { + const [keys, setKeys] = useState([]); + const [secret, setSecret] = useState(''); + usePendingCloudInput(Boolean(secret)); + const [busy, setBusy] = useState(false); + const [error, setError] = useState(''); + const operation = useRef(null); + const refresh = useCallback(async () => { const result = await cloudAgentsApi.call('keys', { application_id: application.id }); setKeys(result.keys); }, [application.id]); + useEffect(() => { let live = true; cloudAgentsApi.call('keys', { application_id: application.id }).then(v => { if (live) setKeys(v.keys); }).catch(e => { if (live) setError(errorText(e)); }); return () => { live = false; }; }, [application.id]); + const act = async (task: () => Promise) => { + setBusy(true); setError(''); + try { await task(); await refresh(); await onUpdated(); } catch(e) { setError(errorText(e)); } finally { setBusy(false); } + }; + return
+

{application.name} · {application.enabled ? '已启用' : '已停用'}

+
+ {keys.map(key =>
{key.prefix}… + {key.revoked_at ? 已撤销 : }
)} + + {secret &&

请保存这份凭据,关闭此页后不再展示。

+ +
+
} + {error &&

{error}

} +
; +} diff --git a/src/pages/CloudAgents/CloudChat.tsx b/src/pages/CloudAgents/CloudChat.tsx new file mode 100644 index 0000000..1ab6a92 --- /dev/null +++ b/src/pages/CloudAgents/CloudChat.tsx @@ -0,0 +1,293 @@ +import { useCallback, useEffect, useRef, useState } from 'react'; +import { MessageSquarePlus, Send, Square } from 'lucide-react'; +import { Button } from '@/components/ui/button'; +import { Textarea } from '@/components/ui/textarea'; +import { cloudAgentsApi } from '@/lib/cloud-agents-api'; +import type { CloudHistory, CloudInterrupt, CloudPrompt, CloudRequest, CloudRun, CloudThread } from '../../../shared/cloud-agents'; +import { CloudFiles } from './CloudFiles'; +import ReactMarkdown, { defaultUrlTransform } from 'react-markdown'; +import remarkGfm from 'remark-gfm'; +import { CloudPendingContext, useCloudPendingState, usePendingCloudInput } from './CloudPending'; + +export const cloudStatus = (status: string) => ({ + idle: '尚未开始', queued: '等待执行', dispatching: '准备执行', dispatched: '已提交', pending: '准备中', + running: '运行中', cancel_requested: '正在停止', cancelled: '已停止', completed: '已完成', + failed: '失败', interrupted: '等待你的确认', submitted: '已提交', rejected: '未执行', +}[status] ?? status); +const active = (status?: string) => Boolean(status && ['queued', 'dispatching', 'dispatched', 'pending', 'running', 'cancel_requested', 'submitted'].includes(status)); +const message = (e: unknown) => e instanceof Error ? e.message : '暂时无法完成操作'; +const object = (v: unknown): Record => v && typeof v === 'object' ? v as Record : {}; + +export function CloudChat({ slug, previewRevision, initialThread }: { slug: string; previewRevision?: number; initialThread?: CloudThread }) { + const [threads, setThreads] = useState([]); + const [thread, setThread] = useState(initialThread ?? null); + const [clientId, setClientId] = useState(() => initialThread?.client_thread_id ?? crypto.randomUUID()); + const [history, setHistory] = useState(null); + const [request, setRequest] = useState(null); + const [run, setRun] = useState(null); + const [query, setQuery] = useState(''); + const [intent, setIntent] = useState(null); + const [busy, setBusy] = useState(false); + const [error, setError] = useState(''); + const [liveText, setLiveText] = useState(''); + const [progress, setProgress] = useState(''); + const [interrupt, setInterrupt] = useState(null); + const [loading, setLoading] = useState(true); + const [leaving, setLeaving] = useState(false); + const [attachmentIds, setAttachmentIds] = useState([]); + const childrenPending = useCloudPendingState(); + usePendingCloudInput(busy || Boolean(intent) || Boolean(query.trim()) || attachmentIds.length > 0 || childrenPending.hasPending); + const generation = useRef(0); + const alive = useRef(true); + const bottom = useRef(null); + const mode = previewRevision === undefined ? 'published' : 'preview'; + const readHistory = useCallback(async (id: string, expectedGeneration = generation.current) => { + let page = await cloudAgentsApi.call('history', { thread_id: id }); + while (page.next_offset !== null) { + const more = await cloudAgentsApi.call('history', { thread_id: id, offset: page.next_offset }); + page = { ...more, messages: [...page.messages, ...more.messages] }; + } + if (!alive.current || expectedGeneration !== generation.current) return; + setHistory(page); setRun(page.run); setInterrupt(page.run?.interrupt ?? null); + if (page.run) setThreads(current => current.map(item => item.thread_id === id + ? { ...item, run_id: page.run!.agent_run_id, status: page.run!.status, unread: false } : item)); + if (page.run) void cloudAgentsApi.call('viewed', { thread_id: id, run_id: page.run.agent_run_id }).catch(() => undefined); + }, []); + useEffect(() => { + alive.current = true; + const g = ++generation.current; + const open = async () => { + try { + if (previewRevision !== undefined) { + if (initialThread) await readHistory(initialThread.thread_id, g); + return; + } + let page = await cloudAgentsApi.call('threads', { slug }); + while (page.next_offset !== null) { + const more = await cloudAgentsApi.call('threads', { slug, offset: page.next_offset }); + page = { ...more, threads: [...page.threads, ...more.threads] }; + } + if (!alive.current || g !== generation.current) return; + const items = page.threads.filter(t => t.mode === mode); + setThreads(items); + const latest = initialThread ?? items[0]; + if (latest) { + setThread(latest); setClientId(latest.client_thread_id ?? latest.thread_id); + await readHistory(latest.thread_id, g); + } + } catch (e) { if (alive.current && g === generation.current) setError(message(e)); } + finally { if (alive.current && g === generation.current) setLoading(false); } + }; + void open(); + return () => { alive.current = false; generation.current++; }; + }, [slug, mode, previewRevision, initialThread, readHistory]); + + useEffect(() => { + if (!request || request.run_id || !active(request.status)) return; + let live = true; + const poll = async () => { + try { + const next = await cloudAgentsApi.call('request', { request_id: request.request_id }); + if (live) { + setRequest(next); + if (next.run_id) { + const acceptedRun = await cloudAgentsApi.call('run', { run_id: next.run_id }); + if (live) setRun(acceptedRun); + } + } + } catch (e) { if (live) setError(message(e)); } + }; + const timer = window.setInterval(() => void poll(), 1800); + void poll(); + return () => { live = false; window.clearInterval(timer); }; + }, [request]); + + const runId = run?.agent_run_id ?? request?.run_id; + const running = active(run?.status ?? request?.status); + useEffect(() => { + if (!runId || !running) return; + let live = true; + let source: EventSource | undefined; + const g = generation.current; + const accept = (event: MessageEvent) => { + if (!live || g !== generation.current) return; + try { + const envelope = object(JSON.parse(event.data)); + const payload = object(envelope.payload); + const chunk = object(payload.chunk); + const targetThread = thread?.thread_id ?? request?.thread_id; + if (envelope.thread_id && targetThread && envelope.thread_id !== targetThread) return; + for (const item of Array.isArray(payload.items) ? payload.items : [chunk]) { + const semantic = object(object(item).stream_event); + if (semantic.type === 'message_delta' && typeof semantic.content === 'string') setLiveText(v => v + semantic.content); + if (semantic.type === 'tool_call' && typeof semantic.name === 'string') setProgress('正在使用 ' + semantic.name); + } + if (event.type === 'interrupt') setInterrupt(chunk as CloudInterrupt); + if (event.type === 'end') { + source?.close(); + setLiveText(''); setProgress(''); + const id = thread?.thread_id ?? request?.thread_id; + if (id) void readHistory(id, g).catch(e => { if (live) setError(message(e)); }); + } + } catch { setError('运行事件读取失败,正在重新读取状态'); } + }; + cloudAgentsApi.events(runId).then(stream => { + if (!live) { stream.close(); return; } + source = stream; + for (const name of ['metadata', 'messages', 'custom', 'interrupt', 'end']) stream.addEventListener(name, accept as EventListener); + stream.onerror = () => { if (live) setProgress('连接中断,正在重连…'); }; + stream.onopen = () => { if (live) setProgress(''); }; + }).catch(e => { if (live) setError(message(e)); }); + // Durable state is also read after an event stream expires or its terminal event was missed. + const timer = window.setInterval(() => { + cloudAgentsApi.call('run', { run_id: runId }).then(next => { + if (!live || g !== generation.current) return; + setRun(next); + if (!active(next.status)) { + source?.close(); setLiveText(''); setProgress(''); setInterrupt(next.interrupt ?? null); + void readHistory(next.thread_id, g).catch(e => { if (live) setError(message(e)); }); + } + }).catch(e => { if (live) setError(message(e)); }); + }, 5000); + return () => { live = false; source?.close(); window.clearInterval(timer); }; + }, [runId, running, thread?.thread_id, request?.thread_id, readHistory]); + useEffect(() => { bottom.current?.scrollIntoView?.({ block: 'end' }); }, [history?.messages.length, liveText, interrupt]); + + const rememberThread = (next: CloudThread) => { + setThread(next); + setThreads(current => [next, ...current.filter(item => item.thread_id !== next.thread_id)]); + }; + const send = async () => { + if (busy || running || (!intent && !query.trim())) return; + const input = intent ?? { + request_id: crypto.randomUUID(), thread_id: thread?.thread_id ?? clientId, query: query.trim(), + attachment_file_ids: attachmentIds, + ...(previewRevision === undefined ? {} : { expected_revision: previewRevision }), + }; + setBusy(true); setIntent(input); setError(''); + const g = generation.current; + try { + const accepted = await cloudAgentsApi.call(previewRevision === undefined ? 'submit' : 'preview', { slug, ...input }); + if (!alive.current || g !== generation.current) return; + setRequest(accepted); setIntent(null); setQuery(''); setLiveText(''); setAttachmentIds([]); + rememberThread({ thread_id: accepted.thread_id, client_thread_id: clientId, agent_slug: slug, title: input.query.slice(0, 80), + mode, updated_at: new Date().toISOString(), run_id: accepted.run_id, status: accepted.status, unread: false }); + await readHistory(accepted.thread_id, g); + } catch (e) { if (alive.current && g === generation.current) setError(message(e)); } + finally { if (alive.current && g === generation.current) setBusy(false); } + }; + const choose = async (next: CloudThread | null) => { + const g = ++generation.current; + setThread(next); setClientId(next?.client_thread_id ?? crypto.randomUUID()); setHistory(null); + setRequest(null); setRun(null); setLiveText(''); setInterrupt(null); setError(''); setQuery(''); + setAttachmentIds([]); + if (next) { setLoading(true); try { await readHistory(next.thread_id, g); } catch(e) { setError(message(e)); } finally { setLoading(false); } } + }; + const cancel = async () => { + setBusy(true); setError(''); + try { + if (request) await cloudAgentsApi.call('cancelRequest', { request_id: request.request_id }); + else if (runId) await cloudAgentsApi.call('cancelRun', { run_id: runId }); + } catch (e) { setError(message(e)); } finally { setBusy(false); } + }; + return
+
+
{previewRevision === undefined ? '正式对话' : '草稿预览 · 修订 ' + previewRevision} +

{run?.version && `版本 ${run.version} · `}本次使用由智能体创建者支付词元点数

+ +
+ {threads.length > 0 && } +
+ {loading ?

正在读取对话…

+ : !history?.messages.length && !liveText &&
告诉它你想完成什么。
关闭客户端后,已提交的任务仍在云端继续。
} + {history?.messages.map(item =>
+

{item.role === 'user' ? '你' : '智能体'}

+
url.startsWith('sandbox:/') ? url.slice(8) : defaultUrlTransform(url)} remarkPlugins={[remarkGfm]} components={{ + a: ({ href, children }) => href && !/^(https?:|mailto:|#)/i.test(href) && thread + ? + : {children}, + }}>{item.content}
+
)} + {liveText &&
{liveText}
} + {run?.error && run.status === 'failed' &&

{run.error.message}

} + {interrupt && runId && { + setInterrupt(null); setLiveText(''); setRun(nextRun); setRequest(null); + }} />} +
+
+ {error &&

{error}

} + { + if (thread) return thread.thread_id; + const created = await cloudAgentsApi.call('createThread', { slug, thread_id: clientId, preview: previewRevision !== undefined, expected_revision: previewRevision }); + rememberThread({ thread_id: created.thread_id, client_thread_id: clientId, agent_slug: slug, title: '新对话', mode, + status: 'idle', updated_at: new Date().toISOString(), run_id: null, unread: false }); + return created.thread_id; + }} /> + {leaving &&
放弃未发送的内容并新建对话? +
} +
{ e.preventDefault(); void send(); }}> +