ai-query对接新版freesun-agent接口的分支
You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 

948 lines
31 KiB

'use strict';
const superagent = require('superagent');
const http = require('http');
const https = require('https');
const fsPromises = require('fs').promises;
const { randomUUID } = require('crypto');
const getRequestUserId = (ctx) => {
return ctx.headers['x-user-id'] || ctx.query.userId || ctx.request.body?.userId;
};
// 获取SuperAgent API基础URL
const getSuperAgentApiBase = (ctx) => {
const baseUrl = ctx.app.fs.config?.superAgent?.baseUrl;
if (!baseUrl) {
throw new Error('未配置 superAgent.baseUrl');
}
return String(baseUrl).replace(/\/+$/, "");
};
// 构建完整的SuperAgent API URL
const buildSuperAgentUrl = (ctx, path) => {
const apiBase = getSuperAgentApiBase(ctx);
return `${apiBase}${path}`;
};
// 构建带参数的API路径
const buildSuperAgentPath = (path, params = {}) => {
let nextPath = path;
Object.keys(params).forEach((key) => {
nextPath = nextPath.replace(`:${key}`, encodeURIComponent(params[key]));
});
return nextPath;
};
// 透传 GET 请求的 query(如 draft、include_content、userId)
const buildSuperAgentUpstreamQuery = (ctx, extraQuery = {}) => {
const merged = { ...(ctx.query || {}), ...(extraQuery || {}) };
const query = {};
Object.entries(merged).forEach(([key, value]) => {
if (value === undefined || value === null || value === '') {
return;
}
if (Array.isArray(value)) {
query[key] = value.map((item) => String(item));
return;
}
query[key] = String(value);
});
return query;
};
const buildSuperAgentDownloadQuery = (ctx) => {
const query = new URLSearchParams();
Object.entries(buildSuperAgentUpstreamQuery(ctx)).forEach(([key, value]) => {
if (Array.isArray(value)) {
value.forEach((item) => query.append(key, item));
return;
}
query.set(key, value);
});
const queryString = query.toString();
return queryString ? `?${queryString}` : '';
};
// 构建请求头
const buildSuperAgentHeaders = (userId, extraHeaders = {}) => {
const headers = { ...extraHeaders };
if (userId) {
headers["X-User-Id"] = String(userId);
}
return headers;
};
const buildRequestOptions = (targetUrl, headers, method = 'POST') => {
const parsedUrl = new URL(targetUrl);
return {
protocol: parsedUrl.protocol,
hostname: parsedUrl.hostname,
port: parsedUrl.port,
path: `${parsedUrl.pathname}${parsedUrl.search}`,
method,
headers,
};
};
const getSuperAgentErrorMessage = (error, fallbackMessage) => {
const responseBody = error?.response?.body;
const responseText = error?.response?.text;
return (
responseBody?.message ||
responseBody?.error ||
responseText ||
error?.message ||
fallbackMessage
);
};
const normalizeTaskSession = (task = {}) => {
const sessionKind = task.task_type === 'conversation' ? 'conversation' : 'task';
return {
...task,
session_id: task.id || task.session_id || '',
type: sessionKind,
session_kind: sessionKind,
reused: Boolean(task.reused),
};
};
const getTaskInteractionSessionId = async (ctx, taskId, userId) => {
const res = await superagent
.get(buildSuperAgentUrl(ctx, buildSuperAgentPath('/api/v1/tasks/:taskId', { taskId })))
.set(buildSuperAgentHeaders(userId));
const interactionSessionId = String(res.body?.created_in_session_id || '').trim();
if (!interactionSessionId) {
throw new Error('任务缺少 created_in_session_id');
}
return interactionSessionId;
};
/**
* 功能:转发 SuperAgent 普通 JSON 接口
* 场景:后端只做鉴权头、路径参数和错误日志封装,上游接口契约仍以 OpenAPI 文档为准
* 注意:
* - 文件下载和 SSE 不走这里,避免误把二进制或长连接当 JSON 处理
* - 默认透传 ctx.query;body 由调用方传入或从 ctx.request.body 提取
*/
const proxySuperAgentJsonRequest = async (ctx, options) => {
const {
method,
path,
params,
query,
body,
successLog,
errorMessage,
} = options;
const startTime = Date.now();
const userId = getRequestUserId(ctx);
const targetPath = buildSuperAgentPath(path, params);
const upstreamQuery = buildSuperAgentUpstreamQuery(ctx, query);
ctx.logger.info(`[superAgent] 开始请求上游接口,method:${method.toUpperCase()},path:${targetPath},userId:${userId || ''}`);
try {
let request = superagent[method](buildSuperAgentUrl(ctx, targetPath))
.set(buildSuperAgentHeaders(userId));
if (Object.keys(upstreamQuery).length) {
request = request.query(upstreamQuery);
}
if (body !== undefined) {
request = request.set('Content-Type', 'application/json').send(body);
}
const res = await request;
ctx.status = 200;
ctx.body = res.body;
ctx.logger.info(`${successLog},耗时:${Date.now() - startTime}ms`);
} catch (error) {
ctx.logger.error(`[superAgent] ${errorMessage}`, error);
ctx.status = error?.status || 400;
ctx.body = {
message: getSuperAgentErrorMessage(error, errorMessage),
};
}
};
/**
* 功能:转发 SuperAgent 文件下载流
* 场景:下载上游生成的标书文件或按路径下载附件
* 注意:
* - 使用原生 http/https 转发,避免把 DOCX 等二进制文件完整加载到内存
* - Koa 需要关闭自动响应接管,否则流式响应头和内容可能被覆盖
*/
const pipeSuperAgentDownloadStream = (ctx, { targetUrl, userId, errorMessage }) => {
return new Promise((resolve, reject) => {
const headers = buildSuperAgentHeaders(userId);
const requestOptions = buildRequestOptions(targetUrl, headers, 'GET');
const client = requestOptions.protocol === 'https:' ? https : http;
const upstreamReq = client.request(requestOptions, (upstreamRes) => {
const statusCode = upstreamRes.statusCode || 500;
const isSuccess = statusCode >= 200 && statusCode < 300;
if (!isSuccess) {
const chunks = [];
upstreamRes.on('data', (chunk) => chunks.push(chunk));
upstreamRes.on('end', () => {
const responseText = Buffer.concat(chunks).toString('utf8');
ctx.status = statusCode;
ctx.body = {
message: responseText || errorMessage,
};
resolve();
});
upstreamRes.on('error', reject);
return;
}
ctx.respond = false;
ctx.res.writeHead(statusCode, {
'Content-Type': upstreamRes.headers['content-type'] || 'application/octet-stream',
'Content-Disposition': upstreamRes.headers['content-disposition'] || 'attachment',
'Cache-Control': upstreamRes.headers['cache-control'] || 'no-cache',
});
upstreamRes.pipe(ctx.res);
upstreamRes.on('end', resolve);
upstreamRes.on('error', reject);
});
upstreamReq.on('error', reject);
ctx.req.on('aborted', () => {
upstreamReq.destroy();
});
upstreamReq.end();
});
};
/**
* 功能:转发 SuperAgent SSE 对话流
* 场景:本系统只做后端代理,真实任务消息流接口由 SuperAgent 服务返回 text/event-stream
* 注意:
* - Koa 默认会接管响应体,SSE 必须设置 ctx.respond = false 后手动写入 ctx.res
* - 不使用 superagent.parse 转发流,避免上游流结束前 Koa 无法稳定返回给浏览器
*/
const pipeSuperAgentChatStream = (ctx, { targetUrl, userId, payload }) => {
return new Promise((resolve, reject) => {
const requestBody = JSON.stringify(payload);
const headers = buildSuperAgentHeaders(userId, {
'Content-Type': 'application/json',
'Accept': 'text/event-stream',
'Content-Length': Buffer.byteLength(requestBody),
});
const requestOptions = buildRequestOptions(targetUrl, headers);
const client = requestOptions.protocol === 'https:' ? https : http;
const upstreamReq = client.request(requestOptions, (upstreamRes) => {
const statusCode = upstreamRes.statusCode || 500;
const isSuccess = statusCode >= 200 && statusCode < 300;
if (!isSuccess) {
const chunks = [];
upstreamRes.on('data', (chunk) => chunks.push(chunk));
upstreamRes.on('end', () => {
const responseText = Buffer.concat(chunks).toString('utf8');
ctx.status = statusCode;
ctx.body = {
message: responseText || 'SuperAgent 对话服务返回异常',
};
resolve();
});
upstreamRes.on('error', reject);
return;
}
ctx.respond = false;
ctx.res.writeHead(statusCode, {
'Content-Type': upstreamRes.headers['content-type'] || 'text/event-stream',
'Cache-Control': upstreamRes.headers['cache-control'] || 'no-cache',
'Connection': 'keep-alive',
'X-Accel-Buffering': 'no',
});
let streamBuffer = "";
let hasMeaningfulStreamEvent = false;
upstreamRes.on('data', (chunk) => {
streamBuffer += chunk.toString('utf8');
const blocks = streamBuffer.split(/\n\n/);
streamBuffer = blocks.pop() || "";
blocks.forEach((block) => {
const trimmedBlock = block.trim();
if (!trimmedBlock) return;
let eventName = "message";
const dataLines = [];
trimmedBlock.split(/\r?\n/).forEach((line) => {
if (line.startsWith("event:")) {
eventName = line.slice(6).trim() || "message";
}
if (line.startsWith("data:")) {
dataLines.push(line.slice(5).trim());
}
});
let eventData = null;
const rawData = dataLines.join("\n");
try {
eventData = rawData ? JSON.parse(rawData) : null;
} catch (error) {
eventData = rawData;
}
if (
eventName === "checkpoint" ||
eventName === "progress" ||
eventName === "checkpoint_done" ||
(eventName === "text" &&
((typeof eventData === "object" && eventData && eventData.text) ||
(typeof eventData === "string" && eventData.trim())))
) {
hasMeaningfulStreamEvent = true;
}
const isEmptyDoneEvent =
eventName === "done" &&
(!eventData || (typeof eventData === "object" && !eventData.text));
if (isEmptyDoneEvent && !hasMeaningfulStreamEvent) {
ctx.logger.log("SuperAgent 对话服务返回空 done 事件");
ctx.res.write(
'event: error\ndata: {"text":"SuperAgent 服务未返回有效内容,请检查模型服务或智能体运行日志"}\n\n'
);
return;
}
ctx.res.write(`${trimmedBlock}\n\n`);
});
});
upstreamRes.on('end', () => {
const lastBlock = streamBuffer.trim();
if (lastBlock) {
ctx.res.write(`${lastBlock}\n\n`);
}
ctx.res.end();
resolve();
});
upstreamRes.on('error', (error) => {
ctx.logger.log(error);
ctx.res.end();
reject(error);
});
});
upstreamReq.on('error', reject);
ctx.req.on('aborted', () => {
upstreamReq.destroy();
});
ctx.res.on('close', () => {
if (!ctx.res.writableEnded) {
upstreamReq.destroy();
}
});
upstreamReq.write(requestBody);
upstreamReq.end();
});
};
/**
* 获取SuperAgent会话列表
* @param {Object} ctx Koa上下文
*/
module.exports.listSessions = async (ctx) => {
const startTime = Date.now();
try {
const userId = getRequestUserId(ctx);
ctx.logger.info(`[superAgent] 开始获取会话列表,userId:${userId || ''}`);
const res = await superagent
.get(buildSuperAgentUrl(ctx, '/api/v1/tasks'))
.set(buildSuperAgentHeaders(userId));
ctx.status = 200;
ctx.body = Array.isArray(res.body) ? res.body.map(normalizeTaskSession) : [];
ctx.logger.info(`[superAgent] 会话列表获取完成,数量:${ctx.body.length},耗时:${Date.now() - startTime}ms`);
} catch (error) {
ctx.logger.error('[superAgent] 获取会话列表失败', error);
ctx.status = error?.status || 400;
ctx.body = {
message: getSuperAgentErrorMessage(error, '获取超级智能体会话失败')
};
}
};
/**
* 创建SuperAgent会话
* @param {Object} ctx Koa上下文
*/
module.exports.createSession = async (ctx) => {
const startTime = Date.now();
try {
const userId = getRequestUserId(ctx);
const requestBody = ctx.request.body || {};
const sessionType =
requestBody.type === 'task' || requestBody.type === 'conversation'
? requestBody.type
: undefined;
const taskType = requestBody.task_type || (sessionType === 'conversation' ? 'conversation' : 'bid');
const title = String(requestBody.title || (sessionType === 'conversation' ? '新对话' : '新任务')).trim();
const folderId = requestBody.folder_id ?? requestBody.folderId;
const upstreamBody = {
session_id: requestBody.session_id || requestBody.sessionId || randomUUID(),
task_type: taskType,
title,
reuse_empty: requestBody.reuse_empty === true,
};
if (folderId !== undefined && folderId !== null && folderId !== '') {
upstreamBody.folder_id = folderId;
}
ctx.logger.info(`[superAgent] 开始创建会话,type:${sessionType || 'default'},taskType:${taskType},userId:${userId || ''}`);
const res = await superagent
.post(buildSuperAgentUrl(ctx, '/api/v1/tasks'))
.set(buildSuperAgentHeaders(userId))
.set('Content-Type', 'application/json')
.send(upstreamBody);
ctx.status = 200;
ctx.body = normalizeTaskSession(res.body);
ctx.logger.info(`[superAgent] 会话创建完成,type:${res.body?.type || sessionType || ''},sessionId:${res.body?.session_id || res.body?.id || ''},耗时:${Date.now() - startTime}ms`);
} catch (error) {
ctx.logger.error('[superAgent] 创建会话失败', error);
ctx.status = error?.status || 400;
ctx.body = {
message: getSuperAgentErrorMessage(error, '创建超级智能体会话失败')
};
}
};
/**
* 获取SuperAgent侧栏树(folders + root sessions)
* @param {Object} ctx Koa上下文
*/
module.exports.getSidebar = async (ctx) => {
await proxySuperAgentJsonRequest(ctx, {
method: 'get',
path: '/api/v1/sidebar',
successLog: `[superAgent] 侧栏树获取完成,realm:${ctx.query?.realm || 'all'}`,
errorMessage: '获取超级智能体侧栏失败',
});
};
/**
* 创建SuperAgent文件夹
* @param {Object} ctx Koa上下文
*/
module.exports.createFolder = async (ctx) => {
const requestBody = ctx.request.body || {};
await proxySuperAgentJsonRequest(ctx, {
method: 'post',
path: '/api/v1/folders',
body: {
realm: requestBody.realm || 'task',
name: requestBody.name || '',
},
successLog: `[superAgent] 文件夹创建完成,realm:${requestBody.realm || 'task'}`,
errorMessage: '创建超级智能体文件夹失败',
});
};
/**
* 更新SuperAgent文件夹
* @param {Object} ctx Koa上下文
*/
module.exports.updateFolder = async (ctx) => {
const folderId = ctx.params.folderId;
if (!folderId) {
ctx.status = 400;
ctx.body = { message: '缺少文件夹 ID' };
return;
}
await proxySuperAgentJsonRequest(ctx, {
method: 'patch',
path: '/api/v1/folders/:folderId',
params: { folderId },
body: ctx.request.body || {},
successLog: `[superAgent] 文件夹更新完成,folderId:${folderId}`,
errorMessage: '更新超级智能体文件夹失败',
});
};
/**
* 删除SuperAgent文件夹
* @param {Object} ctx Koa上下文
*/
module.exports.deleteFolder = async (ctx) => {
const folderId = ctx.params.folderId;
if (!folderId) {
ctx.status = 400;
ctx.body = { message: '缺少文件夹 ID' };
return;
}
await proxySuperAgentJsonRequest(ctx, {
method: 'delete',
path: '/api/v1/folders/:folderId',
params: { folderId },
successLog: `[superAgent] 文件夹删除完成,folderId:${folderId}`,
errorMessage: '删除超级智能体文件夹失败',
});
};
/**
* 更新SuperAgent会话(title / folder_id / sort_order)
* @param {Object} ctx Koa上下文
*/
module.exports.patchSession = async (ctx) => {
const sessionId = ctx.params.sessionId;
if (!sessionId) {
ctx.status = 400;
ctx.body = { message: '缺少会话 ID' };
return;
}
await proxySuperAgentJsonRequest(ctx, {
method: 'patch',
path: '/api/v1/tasks/:sessionId',
params: { sessionId },
body: ctx.request.body || {},
successLog: `[superAgent] 会话更新完成,sessionId:${sessionId}`,
errorMessage: '更新超级智能体会话失败',
});
};
/**
* 获取SuperAgent会话历史消息
* @param {Object} ctx Koa上下文
*/
module.exports.getSessionHistory = async (ctx) => {
try {
const sessionId = ctx.params.sessionId;
const userId = getRequestUserId(ctx);
const res = await superagent
.get(buildSuperAgentUrl(ctx, buildSuperAgentPath('/api/v1/tasks/:sessionId/history', { sessionId })))
.set(buildSuperAgentHeaders(userId));
ctx.status = 200;
ctx.body = res.body;
} catch (error) {
ctx.logger.log(error);
ctx.status = 400;
ctx.body = {
message: error?.message || '获取超级智能体历史消息失败'
};
}
};
/**
* 获取SuperAgent会话实时状态
* @param {Object} ctx Koa上下文
*/
module.exports.getSessionState = async (ctx) => {
const sessionId = ctx.params.sessionId;
if (!sessionId) {
ctx.status = 400;
ctx.body = { message: '缺少会话 ID' };
return;
}
await proxySuperAgentJsonRequest(ctx, {
method: 'get',
path: '/api/v1/tasks/:sessionId/state',
params: { sessionId },
successLog: `[superAgent] 会话状态获取完成,sessionId:${sessionId}`,
errorMessage: '获取超级智能体会话状态失败',
});
};
/**
* 获取SuperAgent会话章节列表
* @param {Object} ctx Koa上下文
*/
module.exports.getSessionSections = async (ctx) => {
const sessionId = ctx.params.sessionId;
if (!sessionId) {
ctx.status = 400;
ctx.body = { message: '缺少会话 ID' };
return;
}
await proxySuperAgentJsonRequest(ctx, {
method: 'get',
path: '/api/sessions/:sessionId/sections',
params: { sessionId },
query: { include_content: ctx.query?.include_content || 'true' },
successLog: `[superAgent] 会话章节列表获取完成,sessionId:${sessionId}`,
errorMessage: '获取超级智能体章节列表失败',
});
};
/**
* 获取SuperAgent会话单个章节
* @param {Object} ctx Koa上下文
*/
module.exports.getSessionSection = async (ctx) => {
const sessionId = ctx.params.sessionId;
const sectionId = ctx.params.sectionId;
if (!sessionId || !sectionId) {
ctx.status = 400;
ctx.body = { message: '缺少会话 ID 或章节 ID' };
return;
}
await proxySuperAgentJsonRequest(ctx, {
method: 'get',
path: '/api/sessions/:sessionId/sections/:sectionId',
params: { sessionId, sectionId },
successLog: `[superAgent] 会话章节内容获取完成,sessionId:${sessionId},sectionId:${sectionId}`,
errorMessage: '获取超级智能体章节内容失败',
});
};
/**
* 更新SuperAgent会话章节内容
* @param {Object} ctx Koa上下文
*/
module.exports.updateSessionSection = async (ctx) => {
const sessionId = ctx.params.sessionId;
const sectionId = ctx.params.sectionId;
const requestBody = ctx.request.body || {};
const content = requestBody.content;
const status = requestBody.status || 'edited';
if (!sessionId || !sectionId) {
ctx.status = 400;
ctx.body = { message: '缺少会话 ID 或章节 ID' };
return;
}
if (typeof content !== 'string') {
ctx.status = 400;
ctx.body = { message: '缺少章节正文内容' };
return;
}
const upstreamBody = { content, status };
if (requestBody.title !== undefined && requestBody.title !== null) {
upstreamBody.title = requestBody.title;
}
if (requestBody.order !== undefined && requestBody.order !== null) {
upstreamBody.order = requestBody.order;
}
await proxySuperAgentJsonRequest(ctx, {
method: 'put',
path: '/api/sessions/:sessionId/sections/:sectionId',
params: { sessionId, sectionId },
body: upstreamBody,
successLog: `[superAgent] 会话章节内容更新完成,sessionId:${sessionId},sectionId:${sectionId}`,
errorMessage: '保存超级智能体章节内容失败',
});
};
/**
* 删除SuperAgent会话
* @param {Object} ctx Koa上下文
*/
/**
* 取消SuperAgent任务执行
* @param {Object} ctx Koa上下文
*/
module.exports.cancelSession = async (ctx) => {
const sessionId = ctx.params.sessionId;
const userId = getRequestUserId(ctx);
return proxySuperAgentJsonRequest(ctx, {
method: 'post',
path: '/api/v1/tasks/:sessionId/cancel',
params: { sessionId },
successLog: `[superAgent] 任务取消完成,sessionId:${sessionId}`,
errorMessage: '取消SuperAgent任务失败',
});
};
module.exports.deleteSession = async (ctx) => {
try {
const sessionId = ctx.params.sessionId;
const userId = getRequestUserId(ctx);
const res = await superagent
.delete(buildSuperAgentUrl(ctx, buildSuperAgentPath('/api/v1/tasks/:sessionId', { sessionId })))
.set(buildSuperAgentHeaders(userId));
ctx.status = 200;
ctx.body = res.body;
} catch (error) {
ctx.logger.log(error);
if (error?.status === 404) {
ctx.status = 200;
ctx.body = {
ok: true,
missing: true,
message: 'SuperAgent 会话已不存在',
};
return;
}
ctx.status = 400;
ctx.body = {
message: error?.message || '删除超级智能体会话失败'
};
}
};
/**
* 上传SuperAgent文件
* @param {Object} ctx Koa上下文
*/
module.exports.uploadFile = async (ctx) => {
const startTime = Date.now();
let uploadFilePath = "";
try {
const sessionId = ctx.request.body?.session_id || ctx.query.sessionId;
const userId = getRequestUserId(ctx);
const file = ctx.file || ctx.request.file || ctx.request.files?.file;
if (!sessionId) {
throw new Error('缺少会话 ID');
}
if (!file) {
throw new Error('请上传文件');
}
if (!file.path) {
throw new Error('上传文件解析失败,请重新选择文件');
}
if (!file.size) {
throw new Error('上传文件为空,请重新选择文件');
}
uploadFilePath = file.path;
const originalName = file.originalname || file.name || file.filename || 'upload.bin';
const interactionSessionId = await getTaskInteractionSessionId(ctx, sessionId, userId);
ctx.logger.info(`[superAgent] 开始上传文件,sessionId:${sessionId},fileName:${originalName},userId:${userId || ''}`);
const res = await superagent
.post(buildSuperAgentUrl(ctx, buildSuperAgentPath('/api/v1/tasks/:sessionId/attachments/upload', { sessionId })))
.set(buildSuperAgentHeaders(userId))
.field('session_id', interactionSessionId)
.attach('file', uploadFilePath, originalName)
.buffer(true)
.parse((upstreamRes, callback) => {
const chunks = [];
upstreamRes.on('data', (chunk) => {
chunks.push(chunk);
});
upstreamRes.on('end', () => {
callback(null, Buffer.concat(chunks));
});
});
const responseText = Buffer.isBuffer(res.body) ? res.body.toString('utf8') : '';
const responseBody = responseText ? JSON.parse(responseText) : {};
ctx.status = 200;
ctx.set('Content-Type', 'application/json; charset=utf-8');
ctx.body = responseBody;
ctx.logger.info(`[superAgent] 文件上传完成,sessionId:${sessionId},storage:${responseBody?.storage || ''},耗时:${Date.now() - startTime}ms`);
} catch (error) {
let errorMessage = error?.message || '上传文件失败';
const responseBody = error?.response?.body;
// 上传接口的返回值里包含中文文件名和 URL,必须按 UTF-8 解析上游原始响应,避免本地转发后乱码。
if (Buffer.isBuffer(responseBody)) {
try {
const responseText = responseBody.toString('utf8');
const parsedResponse = responseText ? JSON.parse(responseText) : {};
errorMessage = parsedResponse.message || parsedResponse.error || errorMessage;
} catch (parseError) {
ctx.logger.warn('[superAgent] 上传失败响应解析失败', parseError);
}
}
ctx.logger.error('[superAgent] 上传文件失败', error);
ctx.status = 400;
ctx.body = {
message: errorMessage
};
} finally {
if (uploadFilePath) {
fsPromises.unlink(uploadFilePath).catch(() => {});
}
}
};
/**
* 发送SuperAgent聊天消息(SSE流式响应)
* @param {Object} ctx Koa上下文
*/
module.exports.sendChat = async (ctx) => {
try {
const requestBody = ctx.request.body || {};
const userId = getRequestUserId(ctx);
const sessionId = String(requestBody.session_id || requestBody.sessionId || '').trim();
const content = String(requestBody.content || requestBody.message || '').trim();
if (!sessionId || !content) {
throw new Error('缺少必要参数');
}
const interactionSessionId = await getTaskInteractionSessionId(ctx, sessionId, userId);
const upstreamPayload = {
session_id: interactionSessionId,
content,
};
if (requestBody.display_content) {
upstreamPayload.display_content = String(requestBody.display_content);
}
if (Array.isArray(requestBody.attachment_ids)) {
upstreamPayload.attachment_ids = requestBody.attachment_ids;
}
await pipeSuperAgentChatStream(ctx, {
targetUrl: buildSuperAgentUrl(ctx, buildSuperAgentPath('/api/v1/tasks/:sessionId/messages/stream', { sessionId })),
userId,
payload: upstreamPayload,
});
} catch (error) {
ctx.logger.log(error);
if (ctx.respond === false || ctx.res.headersSent) {
ctx.res.end();
return;
}
ctx.status = 400;
ctx.body = {
message: error?.message || '发送聊天消息失败'
};
}
};
/**
* 回传checkpoint确认结果
* @param {Object} ctx Koa上下文
*/
module.exports.respondCheckpoint = async (ctx) => {
try {
const userId = getRequestUserId(ctx);
const checkpointId = ctx.request.body?.checkpoint_id;
const sessionId = ctx.request.body?.session_id || ctx.request.body?.sessionId;
const action = ctx.request.body?.action;
const edits = ctx.request.body?.edits || {};
const displayContent = ctx.request.body?.display_content;
if (!checkpointId || !action) {
throw new Error('缺少必要参数');
}
if (!sessionId) {
throw new Error('缺少 session_id');
}
await pipeSuperAgentChatStream(ctx, {
targetUrl: buildSuperAgentUrl(
ctx,
buildSuperAgentPath('/api/v1/tasks/:sessionId/checkpoint/respond', { sessionId })
),
userId,
payload: {
checkpoint_id: checkpointId,
action,
edits,
display_content: displayContent,
},
});
} catch (error) {
ctx.logger.log(error);
if (ctx.respond === false || ctx.res.headersSent) {
ctx.res.end();
return;
}
ctx.status = error?.status || 400;
ctx.body = {
message: getSuperAgentErrorMessage(error, '提交确认结果失败')
};
}
};
/**
* 下载文件
* @param {Object} ctx Koa上下文
*/
module.exports.downloadFile = async (ctx) => {
const startTime = Date.now();
try {
const filePath = ctx.query.path;
const sessionId = ctx.query.sessionId || ctx.query.session_id;
const userId = getRequestUserId(ctx);
if (!filePath) {
throw new Error('缺少文件路径');
}
const query = new URLSearchParams({ path: filePath });
if (sessionId) {
query.set('session_id', sessionId);
}
ctx.logger.info(`[superAgent] 开始下载上游文件,sessionId:${sessionId || ''}`);
await pipeSuperAgentDownloadStream(ctx, {
targetUrl: buildSuperAgentUrl(ctx, `/api/download?${query.toString()}`),
userId,
errorMessage: '下载文件失败',
});
ctx.logger.info(`[superAgent] 上游文件下载转发完成,耗时:${Date.now() - startTime}ms`);
} catch (error) {
ctx.logger.error('[superAgent] 下载文件失败', error);
if (ctx.respond === false || ctx.res.headersSent) {
ctx.res.end();
return;
}
ctx.status = 400;
ctx.body = {
message: error?.message || '下载文件失败'
};
}
};
/**
* 下载当前会话生成的完整DOCX
* @param {Object} ctx Koa上下文
*/
module.exports.downloadSessionFile = async (ctx) => {
const startTime = Date.now();
try {
const sessionId = ctx.params.sessionId;
const userId = getRequestUserId(ctx);
if (!sessionId) {
throw new Error('缺少会话 ID');
}
ctx.logger.info(`[superAgent] 开始下载会话生成文件,sessionId:${sessionId}`);
const targetPath = buildSuperAgentPath('/api/sessions/:sessionId/download', { sessionId });
await pipeSuperAgentDownloadStream(ctx, {
targetUrl: buildSuperAgentUrl(ctx, `${targetPath}${buildSuperAgentDownloadQuery(ctx)}`),
userId,
errorMessage: '下载会话生成文件失败',
});
ctx.logger.info(`[superAgent] 会话生成文件下载转发完成,sessionId:${sessionId},耗时:${Date.now() - startTime}ms`);
} catch (error) {
ctx.logger.error('[superAgent] 下载会话生成文件失败', error);
if (ctx.respond === false || ctx.res.headersSent) {
ctx.res.end();
return;
}
ctx.status = 400;
ctx.body = {
message: error?.message || '下载会话生成文件失败',
};
}
};