一次对话的端到端路径:completions API → 调度 → SSE 流式返回
30 秒导读: 用户在 FastGPT 里发一句话,背后是一个 OpenAI 风格的
/chat/completionsHTTP 请求。这一章跟着这条请求走完全程:鉴权 → 取 app 和历史 → 组装成一次工作流运行 → 边跑边用 SSE 把节点产出推回浏览器 → 收尾把这轮对话补全落库。不深入工作流队列内部 (那是 03-workflow-engine),也不深入 LLM/RAG 节点内部 (04-ai-nodes / 05-knowledge-base)。
1. 这是什么(零基础也能懂)
一句话定义: FastGPT 的对话入口是一个兼容 OpenAI /chat/completions 协议的 HTTP
接口,但它背后跑的不是一次纯 LLM 调用,而是一整张工作流图(一个应用可能串了知识库检索、
判断分支、多个 LLM 节点、工具调用……)。
解决什么问题 / 给谁用: 假设你在 FastGPT 上搭了一个客服机器人(一张工作流图),前端聊天框、 或者第三方系统拿着 API Key,都通过同一个 completions 接口来跟它对话。这个接口要负责 把 「一句用户输入」翻译成「这张图跑一遍」,并且一边跑一边把中间结果流式吐回去——让用户看到 打字机效果、看到"正在检索知识库"这样的状态,而不是干等十几秒。
它对外表现成什么样: 就是一次普通的流式 chat 请求。最小示例(stream: true):
# 示意,非源码:一次典型的 FastGPT 对话请求
curl -N https://your-fastgpt/api/v1/chat/completions \
-H "Authorization: Bearer <api-key>" \
-H "Content-Type: application/json" \
-d '{
"chatId": "abc123", # 同一会话的标识,用来串历史
"stream": true, # 要流式
"detail": true, # 要不要把节点级中间过程也推回来
"messages": [{"role":"user","content":"退货政策是什么?"}]
}'
一句话直觉: 把这个接口想成一个翻译官 + 广播员。翻译官把「HTTP 请求」翻成「工作流能 听懂的运行参数」;广播员守在工作流旁边,节点每吐出一点东西(一段文字、一个节点状态、一次工具 调用),就立刻通过 SSE 播报给客户端。真正干活的"工人"是工作流引擎,接口本身只做编排和转述。
本节到此为止,不碰任何底层代码。记住一件事:接口负责 request→run→stream→save 的编排, 不负责图怎么跑。
2. 顶层全景(它大概怎么转)
这一节给你"大盘":一条请求从进门到落库,经过哪几只手。
2.1 主线一图流
怎么读这张图:从上到下是时间顺序;虚线框是"边跑边发生"的流式旁路,不是最后才做。
HTTP POST /api/v1/chat/completions
│
▼
┌───────────────────────────────────────────────────────┐
│ ① 入口 handler(completions.ts) │
│ 解析 body → 鉴权 → 取 app/历史/版本 → 组装运行参数 │
└───────────────────────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────┐
│ ② 预备一轮(preChatRound) │
│ 占用"生成中"锁 + 预创建 Human/AI 占位记录 │
└───────────────────────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────┐
│ ③ 建 SSE 通道(createWorkflowStreamResponseContext) │
│ 写 SSE header + 心跳 + 拿到 workflowResponseWrite │
└───────────────────────────────────────────────────────┘
│
▼
┌───────────────────────────────────────────────────────┐ ┌───────────────────┐
│ ④ 跑工作流(dispatchWorkFlow) │┈┈┈┈┈┈▶│ SSE 旁路(流式) │
│ tracing + usage 记账 + 队列执行 + 清理 ││ 节点 │ answer / 节点状态 │
│ (队列内部机制见 03) ││ 边跑 │ / 工具 / 交互 … │
└───────────────────────────────────────────────────────┘│ 边推 └───────────────────┘
│ ┘
▼
┌───────────────────────────────────────────────────────┐
│ ⑤ 收尾落库(finalizeChatRound) │
│ 把占位记录补全成最终 Human/AI 内容 + 记账 + 发 [DONE]│
└───────────────────────────────────────────────────────┘
2.2 部件一句话职责
| 部件 | 干什么 | 在哪个文件(符号) |
|---|---|---|
| 入口 handler | 解析请求、鉴权、取数据、组装运行参数、串起全流程 | projects/app/src/pages/api/v1/chat/completions.ts:83 handler |
| 鉴权 | 校验 token/API Key、定位 team/成员/app、算权限 | projects/app/src/service/support/permission/auth/chatCompletion.ts:111 authChatCompletionHeaderRequest |
| 预备一轮 | 占"生成中"锁、预创建 Human/AI 占位 chat item | packages/service/core/chat/utils/prepare.ts:221 preChatRound |
| SSE 通道 | 建流式响应上下文,产出 responseWrite 写函数 | packages/service/core/workflow/utils/streamResponseContext.ts:157 createWorkflowStreamResponseContext |
| 工作流入口 | 记账、tracing、创建队列、跑图、收尾清理 | packages/service/core/workflow/dispatch/index.ts:125 dispatchWorkFlow |
| SSE 写函数 | 按 detail/showNodeStatus 过滤事件并写 SSE | packages/service/core/workflow/dispatch/utils/index.ts:216 getWorkflowResponseWrite |
| 落库 | 把本轮 Human/AI 占位补全成最终内容 | packages/service/core/chat/saveChat.ts:221 finalizeChatRound |
| 节点响应存储 | 把每个节点的详细响应树落库 / 留在内存 | packages/service/core/chat/nodeResponseStorage.ts:425 WorkflowNodeResponseWriter |
2.3 主线走一遍(高层,不进代码)
- 进门:
messages里最后一条 user 消息被当作"这一轮的提问",前面的当历史。 - 验明正身: 分享链接走
authShareChat,其余走authChatCompletionHeaderRequest, 拿到teamId / tmbId / app。 - 凑齐材料: 一次
Promise.all并发取「历史消息、app 最新版本(节点+边+配置)、会话级变量」。 - 占坑:
preChatRound抢下"这个 chatId 正在生成"的锁,并预写一对 Human/AI 占位记录。 - 开广播站: 建 SSE 上下文,拿到
workflowResponseWrite。 - 开跑:
dispatchWorkFlow把节点、边、变量、历史、写函数打包,跑整张图;节点产出 实时经workflowResponseWrite变成 SSE 事件。 - 收尾: 图跑完 → 把占位记录补全成最终内容 → 记账 → 流式场景发
[DONE],非流式场景一次性 返回 JSON。
3. 核心原理(逐个机制,由浅入深)
3.1 入口:一个请求怎么变成"运行参数"
要解决的小问题: HTTP body 里是 messages / chatId / stream / variables …,但工作流引擎
要的是 runtimeNodes / runtimeEdges / histories / query …。入口的核心工作就是这层翻译。
思路: 先鉴权拿到 app,再把 app 存的"静态图"(store 节点/边)转成"可运行的图"(runtime
节点/边),同时把历史和本轮提问对齐。
关键步骤(按代码顺序):
-
解析 + 鉴权。 body 用 zod schema 校验(
parseApiInput+CompletionsPropsSchema)。 分享场景和普通场景分流鉴权:// v1/chat/completions.ts:151 —— 鉴权分流(节选)if (shareId && outLinkUid) {return authShareChat({ shareId, outLinkUid, chatId, ip: originIp, question: startHookText });}return authChatCompletionHeaderRequest({ req, appId, chatId, authProxy, showSkillReferences: true });authChatCompletionHeaderRequest内部用authCert认 token/API Key,再按ReadPermissionVal校验对这个 app 的读权限,返回teamId / tmbId / app / apikey / showCite …(chatCompletion.ts:190的return)。团队级 API Key 还能带authProxy指定 "代表团队内某成员执行",effective tmbId 由resolveChatCompletionEffectiveTmbId决定。 -
取提问 + 取材料。 最后一条 user 消息被
pop出来当userQuestion(插件类型除外, 走serverGetWorkflowToolRunUserQuery);然后一次并发把三样东西取齐:// v1/chat/completions.ts:227 —— 并发取历史 / app 版本 / 会话变量const [{ histories }, { versionId, nodes, edges, chatConfig }, chatDetail] = await Promise.all([getChatItems({ ...chatSource, chatId, offset: 0, limit, field: `obj value memories nodeOutputs` }),getAppLatestVersion(app._id, app),MongoChat.findOne({ ...buildChatSourceQuery(chatSource), chatId }, 'source variableList variables')]); -
静态图 → 运行图。
storeNodes2RuntimeNodes/storeEdges2RuntimeEdges把编辑态的 节点/边转成运行态;入口节点由getWorkflowEntryNodeIds决定(普通对话是 workflowStart, 若上一轮停在交互节点则从交互点续跑)。历史和本轮消息用concatHistories拼接,getLastInteractiveValue检测是否处于"交互待续"态。这几步的数据模型细节见 01-workflow-data-model。
要点: 入口不决定图怎么跑,只负责把"请求上下文"翻译成一份完整的运行参数;真正的调度
交给 dispatchWorkFlow。
3.2 预备一轮:先占坑、先落占位记录
要解 决的小问题: 对话是长耗时操作(可能跑十几秒)。如果同一个 chatId 被并发点了两次、
或者中途崩了,历史记录会乱。FastGPT 的答案是:先占坑。
思路: 在真正跑图之前,preChatRound 做三件事(prepare.ts:221):
- 抢生成锁。
tryStartGenerateChat把这个 chatId 标记为generating;抢不到就直接抛ChatErrEnum.chatIsGenerating(v2 会转成 HTTP 409,v2/chat/completions.ts:576)。 - 预创建占位。
prepareChatRound(prepare.ts:118)严格 create(不 upsert)一对 Human/AI chat item,Human 存真实提问、AI 先存空value: [],两者共用同一个roundDataId, 方便前后端用一轮消息 ID 对齐。 - 返回续跑标志。 用
shouldFinalizePreparedRound / shouldPersistChatRound告诉收尾阶段 该"补全占位"还是"更新交互轮"。
为什么先写占位、最后再补全: 见 3.6——这让"边跑边流式"的中间内容有地方挂靠,也让崩溃时 能把这轮标记成 error 而不是凭空多出半条记录。
一个特例: chatId === 'NO_RECORD_HISTORIES'(NO_RECORD_CHAT_ID)表示"这次运行不落库",
isSkipSaveChatId 会让 prepare/finalize 全部短路跳过。
3.3 建 SSE 通道:广播站怎么开张
要解决的小问题: 流式返回要求 HTTP 响应保持长连接、按 SSE 协议一段段写。谁来建这条通道? 必须由 API 入口显式建——工作流引擎只管跑图,不隐式碰响应协议。
SSE(Server-Sent Events) = 服务器通过一条不关闭的 HTTP 连接,持续往客户端推 event:/data:
文本行。FastGPT 用它实现打字机效果和节点状态提示。
思路: createWorkflowStreamResponseContext(streamResponseContext.ts:157)一次性建好:
- 写 SSE header + 心跳。
initWorkflowSseResponse(streamResponseContext.ts:82)设置Content-Type: text/event-stream、X-Accel-Buffering: no等头,并挂一个 10 秒空 answer 心跳,防止浏览器/代理误判长连接已断(streamResponseContext.ts:127)。 - 产出写函数。 返回
responseWrite(即入口里的workflowResponseWrite),这是后面所有 流式事件的唯一出口。 - 可选断线续传镜像。
getStreamResumeMirror会把原始 SSE chunk 镜像到 Redis,支持 断线后 resume(v1 入口显式enableStreamResume: false关掉,v2 默认开并在结尾flushResume)。
防重要点: dispatchWorkFlow 开头有一道断言——SSE 没初始化就拒绝执行:
// dispatch/index.ts:144 —— 引擎不隐式管响应协议
if (stream && res && !isWorkflowSseResponseInitialized(res)) {
return Promise.reject(new Error('Workflow SSE response must be initialized before dispatchWorkFlow'));
}
这条边界很关键:响应协议归 API 层,图执行归引擎层,两者不越界。
3.4 调度入口:dispatchWorkFlow 是一层"运行时封装"
要解决的小问题: 跑一张图不只是"执行节点",还得记账、上报 tracing、管 abort/停止信号、跑完
清理连接。这些横切关注点由 dispatchWorkFlow 统一封装,把干净的队列算法留给 WorkflowQueue。
思路: dispatchWorkFlow(dispatch/index.ts:125)是队列的外壳,做这几件事:
-
前置校验 + 初始化。 校验文件 URL 域名(
validateFileUrlDomain)、检查团队 AI 积分 (checkTeamAIPoints)、取用户时区/外部变量。 -
usage 记账起账。 并发里创建一条用量记录(续跑则复用上一轮 usageId):
// dispatch/index.ts:176 —— 起一条 usage 记录return createChatUsageRecord({appName: runningAppInfo.name,appId: ..., teamId: runningUserInfo.teamId, tmbId: runningUserInfo.tmbId,source: usageSource});之后队列里每个节点的消耗由
pushChatItemUsage挂到这个 usageId 上(dispatch/index.ts:664, 只有 root runtime 推送,子流程统一上交 root)。 -
停止/中断信号。 进场先
delAgentRuntimeStopSign清掉 Redis 里的旧停止标记 (dispatch/index.ts:203)。运行中如何感知"该停了",v1 和 v2 走两套机制:版本 停止判断 机制 v1 客户端断开连接即停 createClientAbortTracker(监听 req/res)v2 轮询 Redis 停止标记 每 100ms shouldWorkflowStop统一由
checkIsStopping()暴露给队列(dispatch/index.ts:226)。 -
建 nodeResponseWriter + 跑队列。 用
createWorkflowEntryNodeResponseWriter建好节点响应 写入器(3.5),然后在一个带runWithContext的 Promise 里调runWorkflow真正跑图。 -
收尾清理(finally)。 无论成败都:清停止轮询定时器、
clientAbortTracker.cleanup()、 关闭所有 mcpClient 连接、再删一次 Redis 停止标记:// dispatch/index.ts:303 —— 跑完关掉工具调用建立的 MCP 连接Object.values(ctx.mcpClientMemory).forEach((client) => {client.closeConnection();});
边界提醒: WorkflowQueue(dispatch/index.ts:343)本身——并发控制、节点 run/skip/wait
判定、回边/SCC 分组——是下一章 03-workflow-engine 的主题,本章
只讲到"入口怎么把它包起来、喂什么参数(RunWorkflowProps,dispatch/index.ts:319)"。
3.5 流式返回:一个事件枚举 + 一个写函数
要解决的小问题: 工作流里各种节点想推的东西五花八门——一段答案文字、"某节点正在运行"、 一次工具调用、一个需要用户选择的交互卡片。怎么统一推给前端?答案是一套事件枚举 + 一个写函数。
事件枚举 SseResponseEventEnum(packages/global/core/workflow/runtime/constants.ts:3)
列举了所有 SSE 事件类型。挑主线相关的几个:
| 事件 | 含义 | 谁发 |
|---|---|---|
answer | 流式答案文本(打字机动画) | LLM 节点边生成边推 |
fastAnswer | 直接答案文本(不做动画) | 指定输出节点 |
flowNodeStatus | 某节点开始运行的状态提示 | 队列在跑节点前推(dispatch/index.ts:807) |
flowNodeResponse | 节点的详细响应(v2) | 节点跑完推(dispatch/index.ts:996) |
toolCall / toolParams / toolResponse | 工具调用三段 | AI 工具节点(见 04) |
interactive | 需要用户交互(选择/表单/暂停) | 命中交互节点(dispatch/index.ts:1469) |
flowResponses | 全部节点响应汇总(结尾一次性) | 入口收尾 |
chatTitle | 自动生成的会话标题 | 标题生成器 |
写函数 workflowResponseWrite 是所有事件的唯一出口。它由 getWorkflowResponseWrite
(dispatch/utils/index.ts:216)构造,核心是两层过滤:
// dispatch/utils/index.ts:257 —— detail=false 时只放最终答案类事件
const notDetailEvent = { chatTitle: 1, answer: 1, fastAnswer: 1 };
if (!detail && !notDetailEvent[event]) return;
// showNodeStatus=false 时,隐藏节点状态和工具过程(对外 API 常关掉)
const statusEvent = { flowNodeStatus: 1, toolCall: 1, toolParams: 1, toolResponse: 1 };
if (!showNodeStatus && statusEvent[event]) return;
detail=false:客户端只想要最终答案 → 过滤掉一切中间事件。showNodeStatus=false:对外 API / 调试要藏运行细节 → 过滤节点状态和工具参数。
过滤后落到最底层的 responseWrite(packages/service/common/response/index.ts:269),
它就是老老实实往 res 里写 SSE 文本行:
// common/response/index.ts:269 —— 最底层:写一行 SSE
event && Write(`event: ${event}\n`);
Write(`data: ${data}\n\n`);
数据流一句话: 节点 → workflowResponseWrite({event, data}) → 按 detail/showNodeStatus
过滤 → responseWrite 写 event:/data: 行 → 浏览器 EventSource 收到。