feat(agents): 完成云智能体工作台与交互流程
This commit is contained in:
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<typeof app.getPath>[0]) => app.getPath(name));
|
||||
ipcMain.handle('app:take-cloud-agent-link', () => takeCloudAgentRoute());
|
||||
ipcMain.handle('app:quit', () => app.quit());
|
||||
ipcMain.handle('app:relaunch', () => {
|
||||
app.relaunch();
|
||||
|
||||
@@ -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',
|
||||
|
||||
86
electron/services/cloud-agent-operations.ts
Normal file
86
electron/services/cloud-agent-operations.ts
Normal file
@@ -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<string, string>; 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<keyof CloudAgentOperations, OperationSpec> = {
|
||||
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<string, unknown>;
|
||||
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<string, unknown>;
|
||||
return Object.fromEntries(spec.output.filter(key => record[key] !== undefined).map(key => [key, record[key]]));
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -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<string, string> = {
|
||||
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<AbortController, WorksSquareAccountBinding>();
|
||||
private activityTimer: ReturnType<typeof setInterval>;
|
||||
private activityPolling = false;
|
||||
private activitySeen: Map<string, string> | null = null;
|
||||
|
||||
constructor(private readonly fetchImpl: (input: string, init?: RequestInit) => Promise<Response> = 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<unknown> {
|
||||
private async pollActivity(): Promise<void> {
|
||||
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<unknown> {
|
||||
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<CloudUpload | null> {
|
||||
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<unknown> {
|
||||
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<FormData | null> {
|
||||
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<Uint8Array> {
|
||||
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<void>) | 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<unknown> {
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user