重构socket单会话单链接,先握手

This commit is contained in:
2026-08-14 22:22:29 +08:00
parent 7e7a3cbf2d
commit c4c2625714
2 changed files with 252 additions and 115 deletions
+175 -73
View File
@@ -69,7 +69,7 @@
</template>
<script setup lang="ts">
import { ref, reactive, computed, onMounted } from 'vue';
import { ref, reactive, computed, onMounted, onBeforeUnmount } from 'vue';
import { ElMessage, ElMessageBox } from 'element-plus';
import Sidebar from './components/Sidebar.vue';
import MainContent from './components/MainContent.vue';
@@ -77,7 +77,7 @@ import InputBar from './components/InputBar.vue';
import TemplateCompleteDialog from './components/TemplateCompleteDialog.vue';
import { applyHomeFormValues } from './utils/flowDsl';
import { getChatModel } from '/@/api/settings/modelConfigV2';
import { openWsExecute } from './utils/wsExecute';
import { connectSessionSocket, sendAgentStart, sendWorkflowStart, sendCancel } from './utils/wsExecute';
import { parseWsMessage, getDelta, getAnswer, getErrorText, getToolCallName, getToolResultText } from './utils/wsMessage';
import type { ExecutionTreeItem } from '/@/api/settings/creation';
import {
@@ -307,9 +307,19 @@ const getList = async () => {
const selectedWorkflowDetail = ref<any>(null);
const mainContentRef = ref<any>(null);
const sendingSessions = reactive<Record<string, boolean>>({});
// 当前正在执行的 WebSocket(用于停止生成);cancelledWs 记录被主动取消的连接,避免其 onClose 误判为正常完成
const currentWs = ref<WebSocket | null>(null);
const cancelledWs = new Set<WebSocket>();
// ===== 会话级长连接(单活跃连接)=====
// 连接状态/helper 见 formatTime/addMessage 之后统一定义;这里声明本轮 handler 与当前连接路由
interface RoundHandler {
onMessage: (raw: any) => void;
onError: (ev: Event) => void;
onClose: (ev: CloseEvent) => void;
// 终止本轮:保留已生成内容 + flush 打字机 + markStopped + 清 sendingSessions
abort: () => void;
}
// 当前进行中的一轮(连接级 onMessage 只路由给 activeHandlersendingSessions 防重入保证任意时刻最多一轮在生成)
let activeHandler: RoundHandler | null = null;
// 当前活跃的会话级连接(单活跃:切会话时关旧开新)
let wsState: { ws: WebSocket; sid: string; sessionId: string; ready: Promise<WebSocket> } | null = null;
const isHistoryWorkflow = ref(false);
// 生成中:有活跃会话且该会话正处于发送状态(发送按钮切换为停止按钮)
@@ -420,6 +430,101 @@ const addMessage = (msg: ChatMessage) => {
else sessionMessages.value.set(id, [msg]);
};
// ===== 会话级长连接管理 =====
// 关闭当前连接:先置空 wsState,旧 socket 的 onclose 路由被「陈旧守卫」(wsState?.ws !== 旧 ws)丢弃
const closeCurrentSocket = () => {
const st = wsState;
if (!st) return;
wsState = null;
try {
if (st.ws.readyState === WebSocket.OPEN || st.ws.readyState === WebSocket.CONNECTING) st.ws.close();
} catch {
/* 忽略关闭异常 */
}
};
// 确保当前会话有可用连接:sid 匹配且 OPEN 则复用;否则关旧建新(单活跃连接,惰性重连)
const ensureSocket = async (sid: string, sessionId: string): Promise<WebSocket | null> => {
const st = wsState;
if (st && st.sid === sid) {
if (st.ws.readyState === WebSocket.OPEN) return st.ws;
if (st.ws.readyState === WebSocket.CONNECTING) {
try {
return await st.ready;
} catch {
/* 建连失败,走重建 */
}
}
}
closeCurrentSocket();
let session: ReturnType<typeof connectSessionSocket> = null;
session = connectSessionSocket({
sessionId,
onMessage: (raw) => {
if (wsState?.ws === session?.ws) activeHandler?.onMessage(raw);
},
onError: (ev) => {
if (wsState?.ws === session?.ws) activeHandler?.onError(ev);
},
onClose: (ev) => {
// 陈旧连接 / 主动关闭 → 忽略
if (wsState?.ws !== session?.ws) return;
wsState = null;
const h = activeHandler;
activeHandler = null;
if (h) h.onClose(ev);
},
});
if (!session) return null;
wsState = { ws: session.ws, sid, sessionId, ready: session.ready };
try {
return await session.ready;
} catch {
closeCurrentSocket();
return null;
}
};
// 标记「已停止生成」:直接写入指定会话 sid 的消息列表(不能用 addMessage——切会话后 activeHistoryId 已指向新会话)
const markStopped = (sid: string) => {
const list = sessionMessages.value.get(sid);
if (!list) return;
const last = list[list.length - 1];
if (last && !last.isUser) {
last.loading = false;
if (!last.content) last.content = '⚠️ 已停止生成';
} else {
list.push({
id: 'msg-' + Date.now() + '-stop',
content: '⚠️ 已停止生成',
time: formatTime(new Date()),
isUser: false,
});
}
};
// 终止当前轮(停止 / 切会话 / 删除 / 卸载共用):发 cancel 帧(不关闭连接)→ 本轮 abort 定格 → 清空 handler
const teardownActiveRound = () => {
const h = activeHandler;
if (!h) return;
const st = wsState;
if (st && st.ws.readyState === WebSocket.OPEN) {
try {
sendCancel(st.ws, !!selectedWorkflowDetail.value);
} catch {
/* 连接异常时忽略,仅做本地定格 */
}
}
h.abort();
activeHandler = null;
};
// 终止当前会话:终止进行中的轮 + 关闭连接
const teardownSession = () => {
teardownActiveRound();
closeCurrentSocket();
};
const handleSend = async (message: string) => {
// 无活跃会话时自动创建
if (!activeHistoryId.value) {
@@ -567,37 +672,9 @@ const handleDeleteMessage = async (msg: ChatMessage) => {
const handleStopGenerate = () => {
const sid = activeHistoryId.value;
if (!sid || !sendingSessions[sid]) return;
const ws = currentWs.value;
if (ws && ws.readyState === WebSocket.OPEN) {
cancelledWs.add(ws);
try {
const isWorkflow = !!selectedWorkflowDetail.value;
ws.send(JSON.stringify({ type: isWorkflow ? 'workflow_cancel' : 'agent_cancel' }));
} catch {
/* 连接异常时忽略,直接关闭 */
}
ws.close();
}
currentWs.value = null;
// 标记正在生成的 AI 消息:已有内容则保留;无内容则显示停止提示
const list = sessionMessages.value.get(sid);
const last = list && list.length > 0 ? list[list.length - 1] : null;
if (last && !last.isUser) {
last.loading = false;
if (!last.content) last.content = '⚠️ 已停止生成';
} else if (list) {
// 工作流执行模式:无 AI 气泡,追加一条停止提示
addMessage({
id: 'msg-' + Date.now() + '-stop',
content: '⚠️ 已停止生成',
time: formatTime(new Date()),
isUser: false,
});
}
// 清理会话状态(runChat 的 onClose 会再走一次 finishChat,幂等)
const session = historyList.value.find((h) => h.id === sid);
if (session && session.status === 'executing') session.status = 'completed';
delete sendingSessions[sid];
// 会话级长连接:只发 cancel 帧终止本轮(不关闭连接,连接保持复用);
// 「已停止生成」标记与 sendingSessions 清理由本轮 handler 的 abort 统一负责
teardownActiveRound();
};
// ===== 工作流执行:选中工作流 → 表单页 WS 执行 → 完成后切回对话页 =====
@@ -609,7 +686,6 @@ const runWorkflow = async (sid: string, sessionId: string, message: string, mc:
}
let finished = false;
let wsFailed = false;
const finishExec = (success: boolean, errorMsg?: string) => {
if (finished) return;
finished = true;
@@ -678,23 +754,10 @@ const runWorkflow = async (sid: string, sessionId: string, message: string, mc:
nodes: nodeInputParams,
};
// 3. 获取当前会话模型(普通对话用用户设置的会话模型 id
let chatModelId: string | number | undefined;
try {
const chatRes: any = await getChatModel();
chatModelId = chatRes?.data?.modelManage?.id || chatRes?.data?.id || chatRes?.data?.modelId;
} catch {
chatModelId = undefined;
}
// 3. 工作流模型由后端按 flow 配置决定,无需前端取模型 id(普通对话取模型见 runChat
// 4. 打开 WebSocket 执行(消息格式待定:收到完成标志或连接关闭视为完成)
const ws = openWsExecute({
sessionId,
flowId: selectedWorkflowDetail.value.id,
modelId: chatModelId,
question: message || '执行工作流',
systemPrompt: '',
flowContent: updatedFlowContent,
// 4. 会话级长连接:本轮处理进 handler,ensureSocket 复用当前会话连接后 send 工作流启动帧
const handler: RoundHandler = {
onMessage: (raw) => {
// 优先按 ReAct 结构化事件判定完成/失败
const msg = parseWsMessage(raw);
@@ -713,21 +776,35 @@ const runWorkflow = async (sid: string, sessionId: string, message: string, mc:
}
},
onError: () => {
wsFailed = true;
finishExec(false, '执行连接失败,请重试');
},
// 长连接下正常不触发(连接保持),仅异常断开兜底
onClose: () => {
// 主动停止(handleStopGenerate 已发 workflow_cancel 并 close)时跳过,
// 避免被误判为正常完成(工作流完成/失败提示由停止逻辑单独处理)
if (ws && cancelledWs.has(ws)) return;
// 连接关闭兜底视为完成;若已触发 onError 则按失败处理
if (!finished) finishExec(!wsFailed);
if (!finished) finishExec(false, '连接已断开,请重试');
},
});
currentWs.value = ws;
// 停止 / 切会话 / 卸载共用:工作流无 AI 气泡,独立定格(不能复用 finishExec——它会追加 ✅/❌ 消息)
abort: () => {
if (finished) return;
finished = true;
curSession.status = 'completed';
markStopped(sid);
delete sendingSessions[sid];
},
};
// 先挂 handler 再等连接,覆盖建连等待窗口内的停止/切会话
activeHandler = handler;
const ws = await ensureSocket(sid, sessionId);
if (!ws) {
finishExec(false, 'WebSocket 初始化失败,请检查服务地址');
return;
}
if (activeHandler !== handler) return; // 等待期被终止/切走 → 不再发送启动帧
sendWorkflowStart(ws, {
flowId: selectedWorkflowDetail.value.id,
flowContent: updatedFlowContent,
systemPrompt: '',
question: message || '执行工作流',
});
} catch (e: any) {
finishExec(false, e?.message || '执行失败,请重试');
}
@@ -740,7 +817,6 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
addMessage({ id: aiMsgId, content: '', time: formatTime(new Date()), isUser: false, loading: true });
let done = false;
let failed = false;
// 思考/生成过程计时起点(收到首帧时开始)
let thinkStart = 0;
@@ -871,11 +947,8 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
chatModelId = undefined;
}
const ws = openWsExecute({
sessionId,
modelId: chatModelId,
question: message || '',
systemPrompt: '',
// 会话级长连接:本轮消息处理整体进 handler,连接只路由到 activeHandler(打字机/thinking 闭包保留在上面)
const handler: RoundHandler = {
onMessage: (raw) => {
const msg = parseWsMessage(raw);
if (!msg) return;
@@ -948,22 +1021,34 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
// 其余(ack / 未知)忽略
},
onError: () => {
failed = true;
finishChat(false, '对话连接失败,请重试');
},
// 长连接下正常不触发(连接保持),仅异常断开兜底
onClose: () => {
// 连接关闭兜底完成(主动停止时 handleStopGenerate 已把 AI 消息标记为「已停止生成」,
// 此处走成功分支:flush 打字机剩余思考、清理会话状态,且不覆盖已有内容)
if (!done) finishChat(!failed);
if (!done) finishChat(false, '连接已断开,请重试');
},
});
currentWs.value = ws;
// 停止 / 切会话 / 卸载共用:保留已生成内容(finishChat 内部 flush 打字机并清 sendingSessions),无内容则标记已停止
abort: () => {
finishChat(true);
markStopped(sid);
},
};
// 先挂 handler 再等连接,覆盖建连等待窗口内的停止/切会话
activeHandler = handler;
const ws = await ensureSocket(sid, sessionId);
if (!ws) {
finishChat(false, 'WebSocket 初始化失败,请检查服务地址');
return;
}
if (activeHandler !== handler) return; // 等待期被终止/切走 → 不再发送启动帧
sendAgentStart(ws, { modelId: chatModelId, question: message });
};
const handleSelectHistory = async (id: string) => {
// 防重复点击同一条,避免误杀正在生成的轮
if (id === activeHistoryId.value) return;
// 会话级长连接(单活跃):切会话先终止旧轮(保留已生成内容)+ 关闭旧连接
teardownSession();
activeHistoryId.value = id;
activeMenu.value = 'chat';
const session = historyList.value.find((h) => h.id === id);
@@ -1001,6 +1086,10 @@ const handleSelectHistory = async (id: string) => {
if (session?.status === 'completed' || session?.status === 'failed' || (session?.status === 'executing' && !sendingSessions[session.id])) {
session.status = undefined;
}
// 真实会话打开即预建连(fire-and-forget);虚拟会话首次提问才建连
if (session && !session.sessionId.startsWith('virtual_')) {
void ensureSocket(id, session.sessionId);
}
};
// 加载会话内结果(session/get 第一页):workflow+chat 混排,按时间倒序
@@ -1089,6 +1178,8 @@ const handleLoadMore = async (sid: string | null | undefined) => {
const createNewSession = () => {
// 新建会话:终止当前会话的轮与连接(单活跃连接)
teardownSession();
const sessionId = `virtual_${Date.now()}_${Math.random().toString(36).slice(2, 11)}`;
historyList.value.unshift({
id: sessionId,
@@ -1119,6 +1210,12 @@ const handleDeleteHistory = async (id: string) => {
} catch {
return;
}
// 删除的是当前活跃/生成中会话 → 终止该轮并关闭其连接;否则仅关闭对应连接
if (activeHistoryId.value === id) {
teardownSession();
} else if (wsState?.sid === id) {
closeCurrentSocket();
}
const isVirtual = id.startsWith('virtual_');
if (!isVirtual) {
try {
@@ -1149,6 +1246,11 @@ onMounted(() => {
getList();
// 不自动创建 / 选中会话,默认显示占位引导页
});
// 组件卸载:关闭当前会话的长连接,避免泄漏
onBeforeUnmount(() => {
teardownSession();
});
</script>
<style scoped lang="scss">
+77 -42
View File
@@ -1,9 +1,8 @@
// ===== 首页 WebSocket 执行工作流/ai-agent/session/wsExecute=====
// 后端示例:
// WebSocketConnectReq{ SessionId, FlowId, FlowContent, ModelId, Question, SystemPrompt }
// 约定:握手 URL query 带 token/sessionId/flowId/modelId/question/systemPrompt
// 复杂对象 flowContent 在握手后通过 ws.send(JSON) 发送
// 服务端推送消息格式待定:本工具只做连接/发送,消息解析交给调用方。
// ===== 首页 WebSocket 会话级长连接/ai-agent/session/wsExecute=====
// 约定:握手 URL query 只带 sessionId/token+可选 modelId/flowId),
// 建连后不发任何帧;每轮提问/工作流执行通过 send* 系列函数发送启动帧,
// 一轮推完由服务端推送完成事件(answer/error)判定结束,连接保持不关闭。
// 服务端推送消息格式见 wsMessage.ts,本工具只做连接/发送,消息解析交给调用方
import { Session } from '/@/utils/storage';
@@ -13,31 +12,32 @@ const getWsBase = (): string => {
return raw.replace(/^http/, 'ws');
};
export interface WsExecuteOptions {
export interface SessionSocket {
ws: WebSocket;
// onopen 时 resolve(ws);打开前 close/error 时 reject,调用方据此判定建连失败
ready: Promise<WebSocket>;
}
export interface ConnectSessionSocketOptions {
sessionId: string;
flowId?: string | number;
modelId?: string | number;
question?: string;
systemPrompt?: string;
flowContent?: any;
onOpen?: (ws: WebSocket) => void;
onMessage?: (raw: any) => void;
onClose?: (ev: CloseEvent) => void;
onError?: (ev: Event) => void;
flowId?: string | number;
onMessage: (raw: any) => void;
onClose: (ev: CloseEvent) => void;
onError: (ev: Event) => void;
}
/**
* 打开 WebSocket 执行工作流
* 返回 WebSocket 实例(若 URL 构造失败返回 null)。
* 连接成功后自动发送 flowContent(握手后 send JSON
* 建立会话级 WebSocket 长连接
* 握手 query 只带 sessionId/token+可选 modelId/flowId),不做任何自动发送;
* 后续每轮提问由 sendAgentStart / sendWorkflowStart 发送启动帧
* 返回 SessionSocketURL 构造失败返回 null。
*/
export function openWsExecute(opts: WsExecuteOptions): WebSocket | null {
export function connectSessionSocket(opts: ConnectSessionSocketOptions): SessionSocket | null {
const base = getWsBase();
const params: Record<string, string> = { sessionId: opts.sessionId };
if (opts.flowId) params.flowId = String(opts.flowId);
if (opts.modelId) params.modelId = String(opts.modelId);
if (opts.question) params.question = opts.question;
if (opts.systemPrompt) params.systemPrompt = opts.systemPrompt;
if (opts.flowId != null) params.flowId = String(opts.flowId);
if (opts.modelId != null) params.modelId = String(opts.modelId);
const token = Session.get('token');
if (token) params.token = String(token);
@@ -54,24 +54,59 @@ export function openWsExecute(opts: WsExecuteOptions): WebSocket | null {
return null;
}
ws.onopen = () => {
if (opts.flowContent) {
// 工作流执行:连接后发送 flowContentworkflow_exec 的完整 payload 结构待后端确认)
ws.send(JSON.stringify({ flowContent: opts.flowContent, systemPrompt: opts.systemPrompt || '' }));
} else {
// 普通对话:连接后需发送启动消息(type=agent),后端据此开始推流
ws.send(
JSON.stringify({
type: 'agent',
payload: { modelId: opts.modelId != null ? String(opts.modelId) : '', question: opts.question || '' },
})
);
}
opts.onOpen?.(ws);
};
ws.onmessage = (ev) => opts.onMessage?.(ev.data);
ws.onclose = (ev) => opts.onClose?.(ev);
ws.onerror = (ev) => opts.onError?.(ev);
let readyResolve!: (w: WebSocket) => void;
let readyReject!: (e: unknown) => void;
const ready = new Promise<WebSocket>((resolve, reject) => {
readyResolve = resolve;
readyReject = reject;
});
return ws;
// 打开前是否已就绪:用于区分「未连上就关闭」与「正常关闭」
let opened = false;
ws.onopen = () => {
// 长连接:握手成功仅标记就绪,不发送任何启动帧
opened = true;
readyResolve(ws);
};
ws.onmessage = (ev) => opts.onMessage(ev.data);
ws.onclose = (ev) => {
if (!opened) readyReject(ev);
opts.onClose(ev);
};
ws.onerror = (ev) => opts.onError(ev);
return { ws, ready };
}
// ===== 启动帧 =====
/** 普通对话:发送提问(type=agent),后端据此开始推流 */
export function sendAgentStart(ws: WebSocket, p: { modelId?: string | number; question: string }): void {
ws.send(
JSON.stringify({
type: 'agent',
payload: { modelId: p.modelId != null ? String(p.modelId) : '', question: p.question || '' },
})
);
}
/** 工作流执行:发送 flowContent(保持原无 type 结构,向后兼容);flowId 从握手 query 移到帧内 */
export function sendWorkflowStart(
ws: WebSocket,
p: { flowId?: string | number; flowContent: any; systemPrompt?: string; question?: string }
): void {
ws.send(
JSON.stringify({
flowId: p.flowId != null ? String(p.flowId) : '',
question: p.question || '',
flowContent: p.flowContent,
systemPrompt: p.systemPrompt || '',
})
);
}
/** 取消当前轮(不关闭连接,长连接保持复用) */
export function sendCancel(ws: WebSocket, isWorkflow: boolean): void {
ws.send(JSON.stringify({ type: isWorkflow ? 'workflow_cancel' : 'agent_cancel' }));
}