多会话并发执行:每会话独立连接与 handler,切换/新建会话不再取消执行中工作流

重构首页 WebSocket 连接模型:wsConnMap/roundHandlers/runTypeMap 按会话独立管理;执行中保留连接、空闲切换释放;消息按 sid 定向写入;修复后台会话完成时误清当前会话输入状态(resetAll 仅限活跃会话)。

Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
2026-08-20 11:46:40 +08:00
co-authored by Claude
parent 3cddff6df2
commit 20bab0997e
+135 -86
View File
@@ -380,8 +380,8 @@ const handleSessionModelSaved = async (model: { id: string; modelName: string })
const selectedWorkflowDetail = ref<any>(null);
const sendingSessions = reactive<Record<string, boolean>>({});
// ===== 会话级长连接(单活跃连接=====
// 连接状态/helper 见 formatTime/addMessage 之后统一定义;这里声明本轮 handler 与当前连接路由
// ===== 会话级长连接(每会话独立连接,支持多会话并发执行=====
// 连接状态/helper 见 formatTime/addMessage 之后统一定义;这里声明每会话 handler 与连接路由
interface RoundHandler {
onMessage: (raw: any) => void;
onError: (ev: Event) => void;
@@ -389,18 +389,18 @@ interface RoundHandler {
// 终止本轮:保留已生成内容 + 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;
// 当前进行中的一轮类型:true=工作流执行,false=普通对话(sendCancel 协议判断用)
// ref 响应式:InputBar 依据它区分「工作流执行中」与「普通对话生成中」,决定禁用/停止按钮展示
const activeRunIsWorkflow = ref(false);
// 每会话一个进行中的轮(连接级 onMessage 按 sid 路由sendingSessions 防重入保证同一会话任意时刻最多一轮在生成)
const roundHandlers = new Map<string, RoundHandler>();
// 每会话一个会话级连接(切会话不再关旧连接,后台执行中的会话连接保留;空闲会话释放
const wsConnMap: Record<string, { ws: WebSocket; sid: string; sessionId: string; ready: Promise<WebSocket> }> = {};
// 每会话进行中的一轮类型:true=工作流执行,false=普通对话(sendCancel 协议判断用)
// 响应式对象:InputBar 依据它区分「工作流执行中」与「普通对话生成中」,由 isWorkflowRunning 按活跃会话取值
const runTypeMap: Record<string, boolean> = {};
// 生成中:有活跃会话且该会话正处于发送状态(发送按钮切换为停止按钮)
const isGenerating = computed(() => !!activeHistoryId.value && !!sendingSessions[activeHistoryId.value]);
// 工作流执行中:isGenerating 且本轮为工作流 → 输入框整体禁用,停止按钮转移到表单卡片 footer
const isWorkflowRunning = computed(() => isGenerating.value && activeRunIsWorkflow.value);
// 工作流执行中:isGenerating 且当前活跃会话本轮为工作流 → 输入框整体禁用,停止按钮转移到表单卡片 footer
const isWorkflowRunning = computed(() => isGenerating.value && runTypeMap[activeHistoryId.value || ''] === true);
const inputBarRef = ref<any>(null);
const templateDialogVisible = ref(false);
const pendingTemplate = ref<any>(null);
@@ -503,8 +503,20 @@ const claimSessionId = (oldId: string, newId: string) => {
delete sendingSessions[oldId];
}
if (activeHistoryId.value === oldId) activeHistoryId.value = newId;
// 连接状态对齐
if (wsState && wsState.sid === oldId) wsState.sid = newId;
// 连接 / handler / 运行类型 map 对齐(认领后所有以旧 id 为 key 的执行态随迁到新 key)
if (wsConnMap[oldId]) {
wsConnMap[newId] = wsConnMap[oldId];
wsConnMap[newId].sid = newId;
delete wsConnMap[oldId];
}
if (roundHandlers.has(oldId)) {
roundHandlers.set(newId, roundHandlers.get(oldId)!);
roundHandlers.delete(oldId);
}
if (runTypeMap[oldId] !== undefined) {
runTypeMap[newId] = runTypeMap[oldId];
delete runTypeMap[oldId];
}
};
// 认领后刷新会话标题:用后端正式号查一次会话列表,把「新会话 N」更新为后端返回的正式名(fire-and-forget
@@ -645,64 +657,79 @@ 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;
// 按 sid 定向写入消息流(多会话并发:后台会话的完成/失败消息落到其自身列表,不依赖 activeHistoryId
const pushToSession = (sid: string, msg: ChatMessage) => {
const list = sessionMessages.value.get(sid);
if (list) list.push(msg);
else sessionMessages.value.set(sid, [msg]);
};
// ===== 会话级长连接管理(每会话独立连接)=====
// 关闭指定会话的连接:先从 map 移除,旧 socket 的 onclose 路由被「陈旧守卫」(wsConnMap[sid].ws !== 旧 ws)丢弃
const closeSocket = (sid: string) => {
const conn = wsConnMap[sid];
if (!conn) return;
delete wsConnMap[sid];
try {
if (st.ws.readyState === WebSocket.OPEN || st.ws.readyState === WebSocket.CONNECTING) st.ws.close();
if (conn.ws.readyState === WebSocket.OPEN || conn.ws.readyState === WebSocket.CONNECTING) conn.ws.close();
} catch {
/* 忽略关闭异常 */
}
};
// 确保当前会话有可用连接:sid 匹配且 OPEN 则复用;否则关旧建新(单活跃连接,惰性重连
// 确保指定会话有可用连接:sid 已有连接且 OPEN 则复用;CONNECTING 等待;否则为该会话新建(惰性重连,不影响其他会话连接
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) {
const existing = wsConnMap[sid];
if (existing) {
if (existing.ws.readyState === WebSocket.OPEN) return existing.ws;
if (existing.ws.readyState === WebSocket.CONNECTING) {
try {
return await st.ready;
return await existing.ready;
} catch {
/* 建连失败,走重建 */
}
}
}
closeCurrentSocket();
// 会话号可随 ack 认领变化(local → 后端 UUID),holder.sid 跟随认领保持路由一致
const holder = { sid };
let session: ReturnType<typeof connectSessionSocket> = null;
session = connectSessionSocket({
sessionId,
onMessage: (raw) => {
if (wsState?.ws === session?.ws) activeHandler?.onMessage(raw);
const conn = wsConnMap[holder.sid];
if (conn && conn.ws === session?.ws) roundHandlers.get(holder.sid)?.onMessage(raw);
},
onClaim: (uuid) => {
// 首次建连:后端在 ack 里返回正式会话号 → 认领本地临时会话,后续所有操作以正式号为 key
if (wsState?.ws === session?.ws) {
claimSessionId(sid, uuid);
const conn = wsConnMap[holder.sid];
if (conn && conn.ws === session?.ws) {
claimSessionId(holder.sid, uuid);
holder.sid = uuid;
void refreshSessionTitle(uuid);
}
},
onError: (ev) => {
if (wsState?.ws === session?.ws) activeHandler?.onError(ev);
const conn = wsConnMap[holder.sid];
if (conn && conn.ws === session?.ws) roundHandlers.get(holder.sid)?.onError(ev);
},
onClose: (ev) => {
// 陈旧连接 / 主动关闭 → 忽略
if (wsState?.ws !== session?.ws) return;
wsState = null;
const h = activeHandler;
activeHandler = null;
if (h) h.onClose(ev);
// 陈旧连接 / 主动关闭 → 忽略;仅真正断开的当前连接触发会话内 onClose
const conn = wsConnMap[holder.sid];
if (!conn || conn.ws !== session?.ws) return;
delete wsConnMap[holder.sid];
const h = roundHandlers.get(holder.sid);
if (h) {
roundHandlers.delete(holder.sid);
h.onClose(ev);
}
},
});
if (!session) return null;
wsState = { ws: session.ws, sid, sessionId, ready: session.ready };
wsConnMap[holder.sid] = { ws: session.ws, sid: holder.sid, sessionId, ready: session.ready };
try {
return await session.ready;
} catch {
closeCurrentSocket();
delete wsConnMap[holder.sid];
return null;
}
};
@@ -725,26 +752,29 @@ const markStopped = (sid: string) => {
}
};
// 终止当前轮(停止 / 切会话 / 删除 / 卸载共用):发 cancel 帧(不关闭连接)→ 本轮 abort 定格 → 清空 handler
const teardownActiveRound = () => {
const h = activeHandler;
// 终止指定会话的轮(停止 / 删除 / 卸载共用):发 cancel 帧(不关闭连接)→ 本轮 abort 定格 → 清空 handler
// 切会话不再触发终止——多会话并发,后台会话的轮继续执行
const teardownRound = (sid: string) => {
const h = roundHandlers.get(sid);
if (!h) return;
const st = wsState;
if (st && st.ws.readyState === WebSocket.OPEN) {
roundHandlers.delete(sid);
const isWorkflow = runTypeMap[sid] === true;
delete runTypeMap[sid];
const conn = wsConnMap[sid];
if (conn && conn.ws.readyState === WebSocket.OPEN) {
try {
sendCancel(st.ws, activeRunIsWorkflow.value);
sendCancel(conn.ws, isWorkflow);
} catch {
/* 连接异常时忽略,仅做本地定格 */
}
}
h.abort();
activeHandler = null;
};
// 终止当前会话:终止进行中的轮 + 关闭连接
const teardownSession = () => {
teardownActiveRound();
closeCurrentSocket();
// 清理全部会话的轮与连接(组件卸载共用)
const teardownAll = () => {
for (const sid of [...roundHandlers.keys()]) teardownRound(sid);
for (const sid of Object.keys(wsConnMap)) closeSocket(sid);
};
const handleSend = async (message: string) => {
@@ -883,13 +913,13 @@ const handleDeleteMessage = async (msg: ChatMessage) => {
// ===== 停止生成:按执行类型发送对应取消消息并关闭连接,保留已生成内容 =====
// 普通对话走 agent 协议(启动为 {type:'agent'},取消为 {type:'agent_cancel'});
// 工作流走 workflow 协议(取消为 {type:'workflow_cancel'})。
// 判断依据:activeRunIsWorkflow 标记当前轮类型(runWorkflow/runChat 各自设置)。
// 判断依据:runTypeMap[sid] 标记当前会话轮类型(runWorkflow/runChat 各自设置)。
const handleStopGenerate = () => {
const sid = activeHistoryId.value;
if (!sid || !sendingSessions[sid]) return;
// 会话级长连接:只发 cancel 帧终止本轮(不关闭连接,连接保持复用);
// 会话级长连接:只发 cancel 帧终止当前活跃会话本轮(不关闭连接,连接保持复用);
// 「已停止生成」标记与 sendingSessions 清理由本轮 handler 的 abort 统一负责
teardownActiveRound();
teardownRound(sid);
};
// ===== 工作流执行:选中工作流 → 表单页 WS 执行 → 完成后切回对话页 =====
@@ -912,6 +942,11 @@ const runWorkflow = async (
const finishExec = (success: boolean, errorMsg?: string, outputs?: WorkflowOutput[]) => {
if (finished) return;
finished = true;
// 本轮结束:清 handler 与运行类型(正常完成路径不走 teardownRound,主动清理避免残留路由)
roundHandlers.delete(sid);
delete runTypeMap[sid];
// 后台会话(非当前活跃)执行完成 → 释放空闲连接;活跃会话保留连接供继续交互
if (activeHistoryId.value !== sid) closeSocket(sid);
// 定格表单卡片:done/failed + 失败信息;执行结束清进度(避免残留影响重新编辑)
formMsg.formStatus = success ? 'done' : 'failed';
formMsg.formError = success ? undefined : errorMsg || '执行失败';
@@ -930,7 +965,7 @@ const runWorkflow = async (
recordType: roundRecordId ? 'workflow' : undefined,
outputs: outputs || [],
};
addMessage(resultMsg);
pushToSession(sid, resultMsg);
if (success) {
ElMessage.success('✅ 执行完成');
// 刷新工作空间树(查询所有结果)
@@ -943,8 +978,9 @@ const runWorkflow = async (
// 会话在建连时已认领为后端正式号(UUID),执行完成/失败后直接刷新会话结果
loadSessionResults(sid);
delete sendingSessions[sid];
// 执行完成后切回对话页:清空工作流选择,主区域展示消息流
if (success) {
// 执行完成后切回对话页:清空工作流选择,主区域展示消息流
// 仅当当前活跃会话即本轮会话时重置(多会话并发:后台会话完成不干扰当前会话的输入状态)
if (success && activeHistoryId.value === sid) {
inputBarRef.value?.resetAll?.();
}
};
@@ -1008,11 +1044,11 @@ const runWorkflow = async (
onClose: () => {
if (!finished) finishExec(false, '连接已断开,请重试');
},
// 停止 / 切会话 / 卸载共用:工作流无 AI 气泡,独立定格(不能复用 finishExec——它会追加 ✅/❌ 消息)
// 停止 / 删除 / 卸载共用:工作流无 AI 气泡,独立定格(不能复用 finishExec——它会追加 ✅/❌ 消息)
abort: () => {
if (finished) return;
finished = true;
// 取消/切会话:表单卡片定格为已取消(只读、保留已填值),不再追加汇总消息
// 取消/删除:表单卡片定格为已取消(只读、保留已填值),不再追加汇总消息
formMsg.formStatus = 'failed';
formMsg.formError = '执行已取消';
formMsg.formProgress = undefined;
@@ -1020,8 +1056,8 @@ const runWorkflow = async (
delete sendingSessions[sid];
},
};
// 先挂 handler 再等连接,覆盖建连等待窗口内的停止/切会话
activeHandler = handler;
// 先挂 handler 再等连接,覆盖建连等待窗口内的停止
roundHandlers.set(sid, handler);
const ws = await ensureSocket(sid, sessionId);
if (!ws) {
finishExec(false, 'WebSocket 初始化失败,请检查服务地址');
@@ -1030,8 +1066,8 @@ const runWorkflow = async (
// 首次建连已认领后端正式号(ensureSocket 就绪前 ack 已处理):本轮 sid 更新为正式号,
// 后续状态操作(消息/分页/结果)都落到新 key
sid = currentIdOf(sid);
if (activeHandler !== handler) return; // 等待期被终止/切走 → 不再发送启动帧
activeRunIsWorkflow.value = true;
if (roundHandlers.get(sid) !== handler) return; // 等待期被终止 → 不再发送启动帧
runTypeMap[sid] = true;
sendWorkflowStart(ws, {
flowId: detail.id,
flowContent: updatedFlowContent,
@@ -1076,10 +1112,18 @@ const handleFormReEdit = (msg: ChatMessage) => {
msg.formError = undefined;
};
// 工作流执行中取消(卡片 footer「取消执行」):复用 teardownActiveRound —— 发 cancel 帧 + 本轮 abort
// 工作流执行中取消(卡片 footer「取消执行」):复用 teardownRound —— 发 cancel 帧 + 本轮 abort
// abort 把卡片定格为「执行已取消」(保留已填值,可重新编辑并执行)
const handleFormCancel = (_msg: ChatMessage) => {
teardownActiveRound();
const handleFormCancel = (msg: ChatMessage) => {
// 定位卡片所在会话:优先活跃会话,其次按消息列表反查(卡片必在当前活跃会话渲染)
let sid = activeHistoryId.value || '';
for (const [k, list] of sessionMessages.value) {
if (list.includes(msg)) {
sid = k;
break;
}
}
if (sid) teardownRound(sid);
};
// ===== 会话内产出文件卡片操作:预览 / 下载 / 删除 =====
@@ -1126,10 +1170,10 @@ const handleOutputDelete = async (msg: ChatMessage, output: WorkflowOutput) => {
// ===== 普通对话:未选工作流 → 纯问答,AI 回复进消息流 =====
const runChat = async (sid: string, sessionId: string, message: string) => {
activeRunIsWorkflow.value = false;
// AI loading 占位气泡
runTypeMap[sid] = false;
// AI loading 占位气泡(按 sid 定向写入,后台会话的 loading 落在其自身列表)
const aiMsgId = 'msg-' + Date.now() + '-ai';
addMessage({ id: aiMsgId, content: '', time: formatTime(new Date()), isUser: false, loading: true });
pushToSession(sid, { id: aiMsgId, content: '', time: formatTime(new Date()), isUser: false, loading: true });
let done = false;
// 思考/生成过程计时起点(收到首帧时开始)
@@ -1197,6 +1241,11 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
const finishChat = (success: boolean, errorMsg?: string) => {
if (done) return;
done = true;
// 本轮结束:清 handler 与运行类型(正常完成路径不走 teardownRound,主动清理避免残留路由)
roundHandlers.delete(sid);
delete runTypeMap[sid];
// 后台会话(非当前活跃)对话完成 → 释放空闲连接;活跃会话保留连接供继续交互
if (activeHistoryId.value !== sid) closeSocket(sid);
// 结束前清理打字机:成功则把队列剩余思考一次性补全,失败则丢弃未渲染内容
stopTypewriter();
if (success) {
@@ -1214,7 +1263,7 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
list[aiIdx].thinking = '';
list[aiIdx].retryQuestion = message;
} else {
addMessage({ id: 'msg-' + Date.now() + '-fail', content, time: formatTime(new Date()), isUser: false });
pushToSession(sid, { id: 'msg-' + Date.now() + '-fail', content, time: formatTime(new Date()), isUser: false });
}
} else if (list && aiIdx >= 0) {
// 成功兜底:后端未推任何回复内容时,给出明确提示并允许重新生成(避免空白气泡)
@@ -1254,7 +1303,7 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
return;
}
// 会话级长连接:本轮消息处理整体进 handler,连接只路由到 activeHandler(打字机/thinking 闭包保留在上面)
// 会话级长连接:本轮消息处理整体进 handler,连接按 sid 路由到对应会话 handler(打字机/thinking 闭包保留在上面)
const handler: RoundHandler = {
onMessage: (raw) => {
const msg = parseWsMessage(raw);
@@ -1353,14 +1402,14 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
onClose: () => {
if (!done) finishChat(false, '连接已断开,请重试');
},
// 停止 / 切会话 / 卸载共用:保留已生成内容(finishChat 内部 flush 打字机并清 sendingSessions),无内容则标记已停止
// 停止 / 删除 / 卸载共用:保留已生成内容(finishChat 内部 flush 打字机并清 sendingSessions),无内容则标记已停止
abort: () => {
finishChat(true);
markStopped(sid);
},
};
// 先挂 handler 再等连接,覆盖建连等待窗口内的停止/切会话
activeHandler = handler;
// 先挂 handler 再等连接,覆盖建连等待窗口内的停止
roundHandlers.set(sid, handler);
const ws = await ensureSocket(sid, sessionId);
if (!ws) {
finishChat(false, 'WebSocket 初始化失败,请检查服务地址');
@@ -1368,17 +1417,19 @@ const runChat = async (sid: string, sessionId: string, message: string) => {
}
// 首次建连已认领后端正式号(ensureSocket 就绪前 ack 已处理):本轮 sid 更新为正式号
sid = currentIdOf(sid);
if (activeHandler !== handler) return; // 等待期被终止/切走 → 不再发送启动帧
if (roundHandlers.get(sid) !== handler) return; // 等待期被终止 → 不再发送启动帧
sendAgentStart(ws, { modelId: chatModelId, question: message });
};
const handleSelectHistory = async (id: string) => {
// 防重复点击同一条,避免误杀正在生成的轮
if (id === activeHistoryId.value) return;
// 会话级长连接(单活跃):切会话终止旧轮(保留已生成内容)+ 关闭旧连接
teardownSession();
// 会话并发:切会话不再终止旧轮/关闭旧连接,后台执行继续;
// 仅当旧会话当前无执行任务时释放其空闲连接(切回时惰性重连)
const prevId = activeHistoryId.value;
activeHistoryId.value = id;
activeMenu.value = 'chat';
if (prevId && !roundHandlers.has(prevId)) closeSocket(prevId);
// 切换会话不继承工作流选择(历史工作流会话回显时预选上次执行的工作流,可自由切换)
selectedWorkflowDetail.value = null;
const session = historyList.value.find((h) => h.id === id);
@@ -1561,8 +1612,8 @@ const handleLoadMore = async (sid: string | null | undefined) => {
const createNewSession = () => {
// 新建会话:终止当前会话的轮与连接(单活跃连接
teardownSession();
// 新建会话:终止会话的执行(多会话并发,后台继续);旧会话无执行任务时释放空闲连接
const prevId = activeHistoryId.value;
const sessionId = `virtual_${Date.now()}_${Math.random().toString(36).slice(2, 11)}`;
historyList.value.unshift({
id: sessionId,
@@ -1575,6 +1626,7 @@ const createNewSession = () => {
selectedWorkflowDetail.value = null;
inputBarRef.value?.resetAll();
sessionMessages.value.set(sessionId, []); // 清空消息
if (prevId && !roundHandlers.has(prevId)) closeSocket(prevId);
};
const handleCreateHistory = () => {
@@ -1592,12 +1644,9 @@ const handleDeleteHistory = async (id: string) => {
} catch {
return;
}
// 删除的是当前活跃/生成中会话 → 终止该轮并关闭其连接;否则仅关闭对应连接
if (activeHistoryId.value === id) {
teardownSession();
} else if (wsState?.sid === id) {
closeCurrentSocket();
}
// 删除会话:终止其进行中的轮(若有)并关闭其连接(后台执行一并放弃)
teardownRound(id);
closeSocket(id);
// 已认领后端正式号的会话才调后端删除接口;未认领的本地临时会话后端无记录,跳过删除
if (hasBackendSessionId(id)) {
try {
@@ -1631,9 +1680,9 @@ onMounted(() => {
// 不自动创建 / 选中会话,默认显示占位引导页
});
// 组件卸载:关闭当前会话的长连接,避免泄漏
// 组件卸载:清理所有会话的轮与长连接,避免泄漏
onBeforeUnmount(() => {
teardownSession();
teardownAll();
});
</script>