Commit 986b2f1f by Archer Committed by GitHub

refactor(workflow): type workflow SSE events (#7283)

* refactor(workflow): type workflow sse events

* chore: remove workflow sse design doc

* fix(chat): include interactive in non-stream completions

* fix(chat): filter completion interactive payload

* refactor(chat): scope completion response helpers

* docs(chat): align interactive completion payload docs

* fix(chat): use field keyed completion content
parent 3aa6d95f
...@@ -11,19 +11,19 @@ You can find the AppId in your application details URL. ...@@ -11,19 +11,19 @@ You can find the AppId in your application details URL.
# Start a Conversation # Start a Conversation
- Authenticate with an API Key. When calling `chat/completions`, passing `appId` in the request body is recommended. - Authenticate with an API Key. When calling `chat/completions` , passing `appId` in the request body is recommended.
- For OpenAI SDK compatibility, `Authorization: Bearer <apiKey>-<appId>` is also supported. The suffix is only a transport compatibility format and is not stored. - For OpenAI SDK compatibility, `Authorization: Bearer <apiKey>-<appId>` is also supported. The suffix is only a transport compatibility format and is not stored.
- To proxy a team member identity through `authProxy`, the team owner must enable `authProxy` when creating or editing the key. The proxied member must still have permission to access the target app and chat. - To proxy a team member identity through `authProxy` , the team owner must enable `authProxy` when creating or editing the key. The proxied member must still have permission to access the target app and chat.
- Some packages require adding `v1` to the `BaseUrl`. If you get a 404 error, try adding `v1` and retry. - Some packages require adding `v1` to the `BaseUrl` . If you get a 404 error, try adding `v1` and retry.
{/* * 对话现在有`v1`和`v2`两个接口,可以按需使用,v2 自 4.9.4 版本新增,v1 接口同时不再维护 */} {/* * 对话现在有`v1`和`v2`两个接口,可以按需使用,v2 自 4.9.4 版本新增,v1 接口同时不再维护 */}
## Start Chat ## Start Chat
The `v1` chat API is compatible with the `GPT` interface! If you're using the standard `GPT` official API, you can access FastGPT by simply changing the `BaseUrl` and `Authorization`. However, note these rules: The `v1` chat API is compatible with the `GPT` interface! If you're using the standard `GPT` official API, you can access FastGPT by simply changing the `BaseUrl` and `Authorization` . However, note these rules:
- Parameters like `model` and `temperature` are ignored. These values are determined by your workflow configuration. - Parameters like `model` and `temperature` are ignored. These values are determined by your workflow configuration.
- Won't return actual `Token` consumed. If needed, set `detail=true` and manually calculate `tokens` from `responseData`. - Won't return actual `Token` consumed. If needed, set `detail=true` and manually calculate `tokens` from `responseData` .
### Request ### Request
...@@ -100,11 +100,11 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \ ...@@ -100,11 +100,11 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \
- headers.Authorization: Bearer [apikey] - headers.Authorization: Bearer [apikey]
- chatId: string | undefined. - chatId: string | undefined.
- Empty or omitted: FastGPT context is not used, and context is built entirely from `messages`. - Empty or omitted: FastGPT context is not used, and context is built entirely from `messages` .
- Non-empty string: uses `chatId` for the chat, automatically reads messages from the FastGPT session, and uses only the last item in `messages` as the user question. Other messages are ignored. Make sure `chatId` is unique and shorter than 250 characters. - Non-empty string: uses `chatId` for the chat, automatically reads messages from the FastGPT session, and uses only the last item in `messages` as the user question. Other messages are ignored. Make sure `chatId` is unique and shorter than 250 characters.
- messages: Same structure as [GPT chat messages](https://platform.openai.com/docs/api-reference/chat/object). - messages: Same structure as [GPT chat messages](https://platform.openai.com/docs/api-reference/chat/object) .
- responseChatItemId: string | undefined. If provided, FastGPT uses it as the response message ID and stores it in the database. Make sure it is unique under the current `chatId`. - responseChatItemId: string | undefined. If provided, FastGPT uses it as the response message ID and stores it in the database. Make sure it is unique under the current `chatId` .
- detail: Whether to return intermediate values. In `stream` mode, they are separated by `event`; in non-stream mode, they are stored in `responseData`. - detail: Whether to return intermediate values. In `stream` mode, they are separated by `event` ; in non-stream mode, they are stored in `responseData` .
- variables: Module variables. This object replaces `[key]` placeholders in input fields. - variables: Module variables. This object replaces `[key]` placeholders in input fields.
</Tab> </Tab>
...@@ -278,26 +278,33 @@ data: [{"moduleName":"知识库搜索","moduleType":"datasetSearchNode","running ...@@ -278,26 +278,33 @@ data: [{"moduleName":"知识库搜索","moduleType":"datasetSearchNode","running
</Tab> </Tab>
<Tab value="Event Values"> <Tab value="Event Values">
event取值: Event values:
- answer: 返回给客户端的文本(最终会算作回答) - answer: Text returned to the client (counts as the final answer)
- fastAnswer: 指定回复返回给客户端的文本(最终会算作回答) - chatTitle: Chat title generated from the current user question
- toolCall: 执行工具 - fastAnswer: Preset reply text returned to the client (counts as the final answer)
- toolParams: 工具参数 - toolCall: Tool execution
- toolResponse: 工具返回 - toolParams: Tool parameters
- flowNodeStatus: 运行到的节点状态 - toolResponse: Tool response
- flowResponses: 节点完整响应 - flowNodeStatus: Current workflow step status
- updateVariables: 更新变量 - flowResponses: Complete workflow step responses
- error: 报错 - updateVariables: Updated variables
- interactive: Interactive config
- error: Error
</Tab> </Tab>
</Tabs> </Tabs>
### Response ### Response
If your workflow contains interactive nodes, still call this API with `detail=true`. Get interactive node config from `event=interactive` data. For `stream=false`, find `type=interactive` elements in choices. If your workflow contains interactive nodes, still call this API with `detail=true` :
- `stream=true` : read the interactive config from `event=interactive` in `data.interactive` .
- `stream=false` : read the element that contains the `interactive` field from `choices[].message.content` .
The `interactive` payload returned to external callers is display config only. It contains only `type` and `params` ; internal runtime fields such as `entryNodeIds` , `memoryEdges` , `nodeOutputs` , and `nodeResponseId` are not returned. If the workflow internally hits a children / loop / tool wrapper interaction, the API returns the deepest user-facing interaction.
When calling a workflow with interactive nodes, if an interactive node is encountered, it returns immediately with this info: When calling a workflow with interactive steps, if an interaction is encountered, it returns immediately. The examples below show the element inside `choices[].message.content[]` when `stream=false` ; when `stream=true` , `event=interactive` returns `{ "interactive": ... }` as its `data` :
<Tabs items={['User Selection','Form Input']}> <Tabs items={['User Selection','Form Input']}>
<Tab value="User Selection"> <Tab value="User Selection">
...@@ -374,9 +381,9 @@ When calling a workflow with interactive nodes, if an interactive node is encoun ...@@ -374,9 +381,9 @@ When calling a workflow with interactive nodes, if an interactive node is encoun
</Tab> </Tab>
</Tabs> </Tabs>
### Continue Interactive Node ### Continue Interaction
After receiving interactive node info, render your UI to guide user input or selection. Then call this API again to continue the workflow. Use this format: After receiving interactive info, render your UI to guide user input or selection. Then call this API again to continue the workflow. Use this format:
<Tabs items={['User Selection','Form Input']}> <Tabs items={['User Selection','Form Input']}>
<Tab value="User Selection"> <Tab value="User Selection">
...@@ -404,7 +411,7 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \ ...@@ -404,7 +411,7 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \
</Tab> </Tab>
<Tab value="Form Input"> <Tab value="Form Input">
Form input is slightly more complex. Serialize the input as a JSON string for `messages`. Object keys match form keys, values are user inputs. Ensure `chatId` is consistent. Form input is slightly more complex. Serialize the input as a JSON string for `messages` . Object keys match form keys, values are user inputs. Ensure `chatId` is consistent.
```bash ```bash
curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \
...@@ -431,11 +438,11 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \ ...@@ -431,11 +438,11 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \
Plugin API is identical to chat API, with slight parameter differences: Plugin API is identical to chat API, with slight parameter differences:
- 调用插件Type的应用时,接口默认为`detail`模式。 - 调用插件 Type 的应用时,接口默认为 `detail` 模式。
- No need to pass `chatId` since plugins run only once. - No need to pass `chatId` since plugins run only once.
- No need to pass `messages`. - No need to pass `messages` .
- Pass `variables` to represent plugin inputs. - Pass `variables` to represent plugin inputs.
- Get plugin outputs from `pluginData`. - Get plugin outputs from `pluginData` .
### Request ### Request
...@@ -458,8 +465,8 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \ ...@@ -458,8 +465,8 @@ curl --location --request POST 'http://localhost:3000/api/v1/chat/completions' \
<Tabs items={['detail=true,stream=false 响应','detail=true,stream=true 响应','Output Retrieval']}> <Tabs items={['detail=true,stream=false 响应','detail=true,stream=true 响应','Output Retrieval']}>
<Tab value="detail=true, stream=false Response"> <Tab value="detail=true, stream=false Response">
- Find plugin output by locating `moduleType=pluginOutput` in `responseData`. Its `pluginOutput` contains the output. - Find plugin output by locating `moduleType=pluginOutput` in `responseData` . Its `pluginOutput` contains the output.
- Stream output is still available via `choices`. - Stream output is still available via `choices` .
```json ```json
{ {
...@@ -588,7 +595,7 @@ data: [{"nodeId":"fdDgXQ6SYn8v","moduleName":"AI 对话","moduleType":"chatNode" ...@@ -588,7 +595,7 @@ data: [{"nodeId":"fdDgXQ6SYn8v","moduleName":"AI 对话","moduleType":"chatNode"
</Tab> </Tab>
<Tab value="Output Retrieval"> <Tab value="Output Retrieval">
event取值: event 取值:
- answer: 返回给客户端的文本(最终会算作回答) - answer: 返回给客户端的文本(最终会算作回答)
- fastAnswer: 指定回复返回给客户端的文本(最终会算作回答) - fastAnswer: 指定回复返回给客户端的文本(最终会算作回答)
...@@ -605,7 +612,7 @@ event取值: ...@@ -605,7 +612,7 @@ event取值:
# Chat CRUD # Chat CRUD
- The following APIs can be called with any `API Key`. - The following APIs can be called with any `API Key` .
- 4.8.12 and above - 4.8.12 and above
\***\*Important Fields\*\*** \***\*Important Fields\*\***
...@@ -1170,7 +1177,7 @@ curl --location --request POST 'http://localhost:3000/api/core/ai/agent/v2/creat ...@@ -1170,7 +1177,7 @@ curl --location --request POST 'http://localhost:3000/api/core/ai/agent/v2/creat
| 参数名 | 类型 | 必填 | 说明 | | 参数名 | 类型 | 必填 | 说明 |
| ------------- | ------ | ---- | ---------------------------------------------------------- | | ------------- | ------ | ---- | ---------------------------------------------------------- |
| appId | string | ✅ | 应用 Id | | appId | string | ✅ | 应用 ID |
| chatId | string | ✅ | Session ID | | chatId | string | ✅ | Session ID |
| questionGuide | object | | 自定义配置,不传的话,则会根据 appId,取最新发布版本的配置 | | questionGuide | object | | 自定义配置,不传的话,则会根据 appId,取最新发布版本的配置 |
......
...@@ -281,6 +281,7 @@ data: [{"moduleName":"知识库搜索","moduleType":"datasetSearchNode","running ...@@ -281,6 +281,7 @@ data: [{"moduleName":"知识库搜索","moduleType":"datasetSearchNode","running
event 取值: event 取值:
- answer: 返回给客户端的文本(最终会算作回答) - answer: 返回给客户端的文本(最终会算作回答)
- chatTitle: 根据本轮用户问题生成的对话标题
- fastAnswer: 指定回复返回给客户端的文本(最终会算作回答) - fastAnswer: 指定回复返回给客户端的文本(最终会算作回答)
- toolCall: 执行工具 - toolCall: 执行工具
- toolParams: 工具参数 - toolParams: 工具参数
...@@ -288,6 +289,7 @@ event 取值: ...@@ -288,6 +289,7 @@ event 取值:
- flowNodeStatus: 运行到的节点状态 - flowNodeStatus: 运行到的节点状态
- flowResponses: 节点完整响应 - flowResponses: 节点完整响应
- updateVariables: 更新变量 - updateVariables: 更新变量
- interactive: 交互节点配置
- error: 报错 - error: 报错
</Tab> </Tab>
...@@ -295,9 +297,14 @@ event 取值: ...@@ -295,9 +297,14 @@ event 取值:
### 交互节点响应 ### 交互节点响应
如果工作流中包含交互节点,依然是调用该 API 接口,需要设置 `detail=true`,并可以从 `event=interactive` 的数据中获取交互节点的配置信息。如果是 `stream=false`,则可以从 choice 中获取 `type=interactive` 的元素,获取交互节点的选择信息。 如果工作流中包含交互节点,依然是调用该 API 接口,需要设置 `detail=true`
当你调用一个带交互节点的工作流时,如果工作流遇到了交互节点,那么会直接返回,你可以得到下面的信息: - `stream=true`:可从 `event=interactive` 的 `data.interactive` 中获取交互节点配置。
- `stream=false`:可从 `choices[].message.content` 中获取包含 `interactive` 字段的元素。
返回给外部调用方的 `interactive` 是展示配置,只包含 `type` 和 `params`;`entryNodeIds` / `memoryEdges` / `nodeOutputs` / `nodeResponseId` 等内部运行态字段不会返回。若内部命中 children / loop / tool 包装交互,接口会返回最深层面向用户的交互节点。
当你调用一个带交互节点的工作流时,如果工作流遇到了交互节点,那么会直接返回。下面示例展示 `stream=false` 时 `choices[].message.content[]` 中的元素;`stream=true` 时 `event=interactive` 的 `data` 为 `{ "interactive": ... }`:
<Tabs items={['用户选择','表单输入']}> <Tabs items={['用户选择','表单输入']}>
<Tab value="用户选择"> <Tab value="用户选择">
......
---
title: 'V4.15.2 (In Progress)'
description: 'FastGPT V4.15.2 Release Notes'
---
## 🚀 New Features
1. Enterprise verification / company verification.
2. The portal page now supports selecting Agent V2 apps for conversations.
## ⚙️ Improvements
1. Updated the delete confirmation copy when a Skill is not associated with an app.
2. Adapted to the latest WeChat publishing channel SDK.
3. Renamed the plugin status from Offline to Uninstalled.
4. Judge nodes now use a unique ID as the identifier instead of the index, so target branches remain stable when branches are deleted or reordered.
5. Files generated by system tools no longer expire after 1 hour. They are now long-lived and deleted together with the conversation.
## 🐛 Fixes
## 🛠️ Code Improvements
1. Refactored Agent V2 assisted generation / ChatAgentHelper to reuse the dialog.
2. AI request records that contain very long base64/data URLs are truncated before saving to prevent possible stack overflows.
3. Unified SSE event wrapping for stronger type hints.
...@@ -5,11 +5,21 @@ description: 'FastGPT V4.15.2 更新说明' ...@@ -5,11 +5,21 @@ description: 'FastGPT V4.15.2 更新说明'
## 🚀 新增内容 ## 🚀 新增内容
1. 企业认证/公司认证能力。
2. 门户页支持选择 AgentV2 应用进行对话。
## ⚙️ 优化 ## ⚙️ 优化
1. Skill 未关联应用时的删除弹窗文案。 1. Skill 未关联应用时的删除弹窗文案。
2. 适配最新微信发布渠道 sdk。 2. 适配最新微信发布渠道 sdk。
3. 将插件的已下线命名改成已卸载。
4. 判断器节点采用唯一 ID 作为标识,而不是 index,实现删除、排序时,目标分支保持不变。
5. 系统工具生成的文件不会 1 小时过期,改成长期,跟随会话一起删除。
## 🐛 修复 ## 🐛 修复
## 🛠️ 代码优化 ## 🛠️ 代码优化
1. Agent V2 辅助生成/ChatAgentHelper 重构,复用对话框。
2. 保存包含超长 base64/data URL 的 AI 请求记录时可能触发栈溢出,提前进行截断。
3. SSE 事件统一封装,强化类型提示。
{ {
"title": "4.15.x", "title": "4.15.x",
"description": "", "description": "",
"pages": ["4151", "4152", "41500", "41507", "41506", "41505", "41504", "41503", "41502", "41501"] "pages": ["4152", "4151", "41500", "41507", "41506", "41505", "41504", "41503", "41502", "41501"]
} }
{ {
"title": "4.15.x", "title": "4.15.x",
"description": "", "description": "",
"pages": ["4151", "4152", "41500", "41507", "41506", "41505", "41504", "41503", "41502", "41501"] "pages": ["4152","4151", "41500", "41507", "41506", "41505", "41504", "41503", "41502", "41501"]
} }
...@@ -161,6 +161,7 @@ description: FastGPT Toc ...@@ -161,6 +161,7 @@ description: FastGPT Toc
- [/en/self-host/upgrading/4-15/41506](/en/self-host/upgrading/4-15/41506) - [/en/self-host/upgrading/4-15/41506](/en/self-host/upgrading/4-15/41506)
- [/en/self-host/upgrading/4-15/41507](/en/self-host/upgrading/4-15/41507) - [/en/self-host/upgrading/4-15/41507](/en/self-host/upgrading/4-15/41507)
- [/en/self-host/upgrading/4-15/4151](/en/self-host/upgrading/4-15/4151) - [/en/self-host/upgrading/4-15/4151](/en/self-host/upgrading/4-15/4151)
- [/en/self-host/upgrading/4-15/4152](/en/self-host/upgrading/4-15/4152)
- [/en/self-host/upgrading/outdated/40](/en/self-host/upgrading/outdated/40) - [/en/self-host/upgrading/outdated/40](/en/self-host/upgrading/outdated/40)
- [/en/self-host/upgrading/outdated/41](/en/self-host/upgrading/outdated/41) - [/en/self-host/upgrading/outdated/41](/en/self-host/upgrading/outdated/41)
- [/en/self-host/upgrading/outdated/4100](/en/self-host/upgrading/outdated/4100) - [/en/self-host/upgrading/outdated/4100](/en/self-host/upgrading/outdated/4100)
......
...@@ -152,7 +152,7 @@ ...@@ -152,7 +152,7 @@
"content/openapi/app.en.mdx": "2026-05-29T19:31:16+08:00", "content/openapi/app.en.mdx": "2026-05-29T19:31:16+08:00",
"content/openapi/app.mdx": "2026-05-29T19:31:16+08:00", "content/openapi/app.mdx": "2026-05-29T19:31:16+08:00",
"content/openapi/chat.en.mdx": "2026-06-23T13:54:06+08:00", "content/openapi/chat.en.mdx": "2026-06-23T13:54:06+08:00",
"content/openapi/chat.mdx": "2026-06-23T13:54:06+08:00", "content/openapi/chat.mdx": "2026-07-08T22:07:36+08:00",
"content/openapi/dataset.en.mdx": "2026-05-29T19:31:16+08:00", "content/openapi/dataset.en.mdx": "2026-05-29T19:31:16+08:00",
"content/openapi/dataset.mdx": "2026-05-29T19:31:16+08:00", "content/openapi/dataset.mdx": "2026-05-29T19:31:16+08:00",
"content/openapi/index.en.mdx": "2026-04-26T21:08:47+08:00", "content/openapi/index.en.mdx": "2026-04-26T21:08:47+08:00",
...@@ -167,8 +167,8 @@ ...@@ -167,8 +167,8 @@
"content/plugin/model-presets.mdx": "2026-06-04T16:10:15+08:00", "content/plugin/model-presets.mdx": "2026-06-04T16:10:15+08:00",
"content/plugin/system-tool-development.en.mdx": "2026-07-02T11:54:55+08:00", "content/plugin/system-tool-development.en.mdx": "2026-07-02T11:54:55+08:00",
"content/plugin/system-tool-development.mdx": "2026-07-02T11:54:55+08:00", "content/plugin/system-tool-development.mdx": "2026-07-02T11:54:55+08:00",
"content/self-host/config/env.en.mdx": "2026-07-02T15:38:53+08:00", "content/self-host/config/env.en.mdx": "2026-07-07T21:14:51+08:00",
"content/self-host/config/env.mdx": "2026-07-02T15:38:53+08:00", "content/self-host/config/env.mdx": "2026-07-07T21:14:51+08:00",
"content/self-host/config/model/intro.en.mdx": "2026-06-04T16:10:15+08:00", "content/self-host/config/model/intro.en.mdx": "2026-06-04T16:10:15+08:00",
"content/self-host/config/model/intro.mdx": "2026-06-04T16:10:15+08:00", "content/self-host/config/model/intro.mdx": "2026-06-04T16:10:15+08:00",
"content/self-host/config/model/minimax.en.mdx": "2026-06-03T10:40:17+08:00", "content/self-host/config/model/minimax.en.mdx": "2026-06-03T10:40:17+08:00",
...@@ -281,6 +281,10 @@ ...@@ -281,6 +281,10 @@
"content/self-host/upgrading/4-14/41426.mdx": "2026-06-24T23:05:29+08:00", "content/self-host/upgrading/4-14/41426.mdx": "2026-06-24T23:05:29+08:00",
"content/self-host/upgrading/4-14/41427.en.mdx": "2026-07-01T17:22:47+08:00", "content/self-host/upgrading/4-14/41427.en.mdx": "2026-07-01T17:22:47+08:00",
"content/self-host/upgrading/4-14/41427.mdx": "2026-07-01T17:22:47+08:00", "content/self-host/upgrading/4-14/41427.mdx": "2026-07-01T17:22:47+08:00",
"content/self-host/upgrading/4-14/41428.en.mdx": "2026-07-06T22:53:13+08:00",
"content/self-host/upgrading/4-14/41428.mdx": "2026-07-06T22:53:13+08:00",
"content/self-host/upgrading/4-14/41429.en.mdx": "2026-07-06T22:53:13+08:00",
"content/self-host/upgrading/4-14/41429.mdx": "2026-07-06T22:53:13+08:00",
"content/self-host/upgrading/4-14/4143.en.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/4-14/4143.en.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/4-14/4143.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/4-14/4143.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/4-14/4144.en.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/4-14/4144.en.mdx": "2026-04-26T21:08:47+08:00",
...@@ -314,8 +318,10 @@ ...@@ -314,8 +318,10 @@
"content/self-host/upgrading/4-15/41506.mdx": "2026-07-01T12:13:58+08:00", "content/self-host/upgrading/4-15/41506.mdx": "2026-07-01T12:13:58+08:00",
"content/self-host/upgrading/4-15/41507.en.mdx": "2026-06-30T17:31:43+08:00", "content/self-host/upgrading/4-15/41507.en.mdx": "2026-06-30T17:31:43+08:00",
"content/self-host/upgrading/4-15/41507.mdx": "2026-06-30T17:31:43+08:00", "content/self-host/upgrading/4-15/41507.mdx": "2026-06-30T17:31:43+08:00",
"content/self-host/upgrading/4-15/4151.en.mdx": "2026-07-04T00:03:29+08:00", "content/self-host/upgrading/4-15/4151.en.mdx": "2026-07-07T21:14:28+08:00",
"content/self-host/upgrading/4-15/4151.mdx": "2026-07-04T00:03:29+08:00", "content/self-host/upgrading/4-15/4151.mdx": "2026-07-07T21:14:28+08:00",
"content/self-host/upgrading/4-15/4152.en.mdx": "2026-07-08T21:29:17+08:00",
"content/self-host/upgrading/4-15/4152.mdx": "2026-07-08T21:29:17+08:00",
"content/self-host/upgrading/outdated/40.en.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/outdated/40.en.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/outdated/40.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/outdated/40.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/outdated/41.en.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/outdated/41.en.mdx": "2026-04-26T21:08:47+08:00",
...@@ -456,6 +462,6 @@ ...@@ -456,6 +462,6 @@
"content/self-host/upgrading/outdated/499.mdx": "2026-05-07T15:06:40+08:00", "content/self-host/upgrading/outdated/499.mdx": "2026-05-07T15:06:40+08:00",
"content/self-host/upgrading/upgrade-intruction.en.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/upgrade-intruction.en.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/upgrade-intruction.mdx": "2026-04-26T21:08:47+08:00", "content/self-host/upgrading/upgrade-intruction.mdx": "2026-04-26T21:08:47+08:00",
"content/toc.en.mdx": "2026-07-01T17:22:47+08:00", "content/toc.en.mdx": "2026-07-08T21:29:17+08:00",
"content/toc.mdx": "2026-07-01T17:22:47+08:00" "content/toc.mdx": "2026-07-06T22:53:13+08:00"
} }
\ No newline at end of file
import type {
ChatHistoryItemResType,
SandboxStatusItemType,
SkillModuleResponseItemType,
ToolModuleResponseItemType
} from '../../chat/type';
import type { AgentPlanStatusType, AgentPlanType } from '../../ai/agent/type';
import type { WorkflowInteractiveResponseType } from '../template/system/interactive/type';
import { SseResponseEventEnum } from './constants';
import { textAdaptGptResponse } from './utils';
export type WorkflowAnswerChunk = ReturnType<typeof textAdaptGptResponse>;
export type WorkflowToolDeltaType = Pick<ToolModuleResponseItemType, 'id'> &
Partial<Omit<ToolModuleResponseItemType, 'id'>>;
export type WorkflowSsePayloadMap = {
[SseResponseEventEnum.error]: Record<string, any>;
[SseResponseEventEnum.workflowDuration]: { durationSeconds: number };
[SseResponseEventEnum.chatTitle]: { title: string };
[SseResponseEventEnum.answer]: WorkflowAnswerChunk;
[SseResponseEventEnum.fastAnswer]: WorkflowAnswerChunk;
[SseResponseEventEnum.flowNodeStatus]: {
status: 'running';
name: string;
};
[SseResponseEventEnum.flowNodeResponse]: ChatHistoryItemResType;
[SseResponseEventEnum.toolCall]: { tool: ToolModuleResponseItemType };
[SseResponseEventEnum.toolParams]: { tool: WorkflowToolDeltaType };
[SseResponseEventEnum.toolResponse]: { tool: WorkflowToolDeltaType };
[SseResponseEventEnum.flowResponses]: Record<string, any>;
[SseResponseEventEnum.updateVariables]: Record<string, any>;
[SseResponseEventEnum.interactive]: { interactive: WorkflowInteractiveResponseType };
[SseResponseEventEnum.plan]: { plan: AgentPlanType };
[SseResponseEventEnum.planStatus]: { planStatus: AgentPlanStatusType };
[SseResponseEventEnum.sandboxStatus]: SandboxStatusItemType;
[SseResponseEventEnum.skillCall]: { skill: SkillModuleResponseItemType };
};
export type WorkflowTypedSseEvent<
Event extends keyof WorkflowSsePayloadMap = keyof WorkflowSsePayloadMap
> = {
id?: string;
event: Event;
data: WorkflowSsePayloadMap[Event];
};
export type WorkflowRawSseEvent = {
id?: string;
event?: SseResponseEventEnum;
data: string;
};
export type WorkflowResponseItemType = WorkflowTypedSseEvent | WorkflowRawSseEvent;
export type WorkflowResponseType = (event: WorkflowResponseItemType) => void;
type AnswerEventParams = {
text?: string | null;
reasoningContent?: string | null;
finishReason?: null | 'stop';
event?: SseResponseEventEnum.answer | SseResponseEventEnum.fastAnswer;
id?: string;
};
const answerEvent = ({
text,
reasoningContent,
finishReason,
event = SseResponseEventEnum.answer,
id
}: AnswerEventParams): WorkflowTypedSseEvent<
SseResponseEventEnum.answer | SseResponseEventEnum.fastAnswer
> => ({
...(id && { id }),
event,
data: textAdaptGptResponse({
text,
reasoning_content: reasoningContent,
finish_reason: finishReason
})
});
export const workflowSseEvent = {
/**
* 输出标准答案文本增量,对应主回答气泡的逐字追加内容。
*/
answerDelta(text: string, id?: string): WorkflowTypedSseEvent<SseResponseEventEnum.answer> {
return answerEvent({ text, id }) as WorkflowTypedSseEvent<SseResponseEventEnum.answer>;
},
/**
* 输出快速回答文本增量,用于 workflow 正式执行结果之外的即时提示内容。
*/
fastAnswerDelta(
text: string,
id?: string
): WorkflowTypedSseEvent<SseResponseEventEnum.fastAnswer> {
return answerEvent({
text,
id,
event: SseResponseEventEnum.fastAnswer
}) as WorkflowTypedSseEvent<SseResponseEventEnum.fastAnswer>;
},
/**
* 输出模型思考内容增量,复用 answer 事件承载 reasoning_content。
*/
reasoningDelta(text: string, id?: string): WorkflowTypedSseEvent<SseResponseEventEnum.answer> {
return answerEvent({
reasoningContent: text,
id
}) as WorkflowTypedSseEvent<SseResponseEventEnum.answer>;
},
/**
* 输出回答结束标记,通知前端当前 answer 或 fastAnswer 文本流已经停止。
*/
answerStop(
event:
| SseResponseEventEnum.answer
| SseResponseEventEnum.fastAnswer = SseResponseEventEnum.answer
): WorkflowTypedSseEvent<SseResponseEventEnum.answer | SseResponseEventEnum.fastAnswer> {
return answerEvent({
text: null,
finishReason: 'stop',
event
});
},
/**
* 输出 SSE 完成标记,data 固定为 [DONE] 以兼容现有前端和 OpenAI 风格流。
*/
done(event?: SseResponseEventEnum): WorkflowRawSseEvent {
return {
event,
data: '[DONE]'
};
},
/**
* 输出已经序列化好的底层事件,用于动态事件名或必须保持旧 wire 格式的兼容场景。
*/
raw({ event, data, id }: WorkflowRawSseEvent): WorkflowRawSseEvent {
return {
...(id && { id }),
event,
data
};
},
/**
* 输出工作流节点运行状态,前端据此展示当前正在执行的节点名称。
*/
flowNodeStatus(name: string): WorkflowTypedSseEvent<SseResponseEventEnum.flowNodeStatus> {
return {
event: SseResponseEventEnum.flowNodeStatus,
data: {
status: 'running',
name
}
};
},
/**
* 输出单个节点响应结果,承载最终会进入 chat history 的节点级响应数据。
*/
flowNodeResponse(
nodeResponse: ChatHistoryItemResType
): WorkflowTypedSseEvent<SseResponseEventEnum.flowNodeResponse> {
return {
event: SseResponseEventEnum.flowNodeResponse,
data: nodeResponse
};
},
/**
* 输出工具调用开始事件,需要包含完整工具展示信息用于初始化工具运行卡片。
*/
toolCall(tool: ToolModuleResponseItemType): WorkflowTypedSseEvent<SseResponseEventEnum.toolCall> {
return {
id: tool.id,
event: SseResponseEventEnum.toolCall,
data: { tool }
};
},
/**
* 输出工具参数增量,只要求携带工具 id 和本次变化字段,前端按 id 合并。
*/
toolParams(tool: WorkflowToolDeltaType): WorkflowTypedSseEvent<SseResponseEventEnum.toolParams> {
return {
id: tool.id,
event: SseResponseEventEnum.toolParams,
data: { tool }
};
},
/**
* 输出工具响应增量,只要求携带工具 id 和本次响应片段,前端按 id 合并。
*/
toolResponse(
tool: WorkflowToolDeltaType
): WorkflowTypedSseEvent<SseResponseEventEnum.toolResponse> {
return {
id: tool.id,
event: SseResponseEventEnum.toolResponse,
data: { tool }
};
},
/**
* 输出完整 workflow 响应集合,通常用于请求结束前同步最终节点响应列表。
*/
flowResponses(
data: Record<string, any>
): WorkflowTypedSseEvent<SseResponseEventEnum.flowResponses> {
return {
event: SseResponseEventEnum.flowResponses,
data
};
},
/**
* 输出变量更新结果,前端据此刷新 workflow 运行后产生的变量值。
*/
updateVariables(
variables: Record<string, any>
): WorkflowTypedSseEvent<SseResponseEventEnum.updateVariables> {
return {
event: SseResponseEventEnum.updateVariables,
data: variables
};
},
/**
* 输出交互节点配置,前端据此渲染需要用户继续操作的交互表单。
*/
interactive(
interactive: WorkflowInteractiveResponseType
): WorkflowTypedSseEvent<SseResponseEventEnum.interactive> {
return {
event: SseResponseEventEnum.interactive,
data: { interactive }
};
},
/**
* 输出工作流总耗时,供前端或调试界面展示本轮执行时长。
*/
workflowDuration(
durationSeconds: number
): WorkflowTypedSseEvent<SseResponseEventEnum.workflowDuration> {
return {
event: SseResponseEventEnum.workflowDuration,
data: { durationSeconds }
};
},
/**
* 输出 Agent 计划内容,使用固定 id 时前端会替换同一条计划展示。
*/
plan(plan: AgentPlanType, id?: string): WorkflowTypedSseEvent<SseResponseEventEnum.plan> {
return {
...(id && { id }),
event: SseResponseEventEnum.plan,
data: { plan }
};
},
/**
* 输出 Agent 计划状态,通常用于标记计划生成、运行或完成状态。
*/
planStatus(
planStatus: AgentPlanStatusType,
id?: string
): WorkflowTypedSseEvent<SseResponseEventEnum.planStatus> {
return {
...(id && { id }),
event: SseResponseEventEnum.planStatus,
data: { planStatus }
};
},
/**
* 输出沙箱执行状态,前端据此展示代码沙箱的运行阶段和结果。
*/
sandboxStatus(
sandboxStatus: SandboxStatusItemType
): WorkflowTypedSseEvent<SseResponseEventEnum.sandboxStatus> {
return {
event: SseResponseEventEnum.sandboxStatus,
data: sandboxStatus
};
},
/**
* 输出技能调用信息,用于在聊天响应中展示 skill 的执行过程。
*/
skillCall(
skill: SkillModuleResponseItemType
): WorkflowTypedSseEvent<SseResponseEventEnum.skillCall> {
return {
event: SseResponseEventEnum.skillCall,
data: { skill }
};
},
/**
* 输出自动生成的聊天标题,前端收到后更新当前会话标题。
*/
chatTitle(title: string): WorkflowTypedSseEvent<SseResponseEventEnum.chatTitle> {
return {
event: SseResponseEventEnum.chatTitle,
data: { title }
};
},
/**
* 输出 workflow 业务错误对象;已序列化错误字符串应使用 raw 保持原格式。
*/
error(data: Record<string, any>): WorkflowTypedSseEvent<SseResponseEventEnum.error> {
return {
event: SseResponseEventEnum.error,
data
};
}
};
...@@ -15,7 +15,7 @@ import type { NextApiResponse } from 'next'; ...@@ -15,7 +15,7 @@ import type { NextApiResponse } from 'next';
import type { AppSchemaType } from '../../app/type'; import type { AppSchemaType } from '../../app/type';
import type { RuntimeEdgeItemType } from '../type/edge'; import type { RuntimeEdgeItemType } from '../type/edge';
import { ReadFileNodeResponseSchema } from '../template/system/readFiles/type'; import { ReadFileNodeResponseSchema } from '../template/system/readFiles/type';
import type { WorkflowResponseType } from '../../../../service/core/workflow/dispatch/type'; import type { WorkflowResponseType } from './sse';
import type { AiChatQuoteRoleType } from '../template/system/aiChat/type'; import type { AiChatQuoteRoleType } from '../template/system/aiChat/type';
import type { OpenaiAccountType } from '../../../support/user/team/type'; import type { OpenaiAccountType } from '../../../support/user/team/type';
import { CompletionFinishReasonSchema } from '../../ai/llm/type'; import { CompletionFinishReasonSchema } from '../../ai/llm/type';
......
...@@ -105,7 +105,7 @@ const ChatCompletionResponseMessageSchema = z.object({ ...@@ -105,7 +105,7 @@ const ChatCompletionResponseMessageSchema = z.object({
role: z.literal('assistant').meta({ description: '消息角色' }), role: z.literal('assistant').meta({ description: '消息角色' }),
content: z.any().meta({ content: z.any().meta({
description: description:
'消息内容。普通对话为字符串;detail=true 或工作流命中交互节点时,可能为带 type 字段的对象数组(type 取值: text / interactive / tool / file / reasoning)' '消息内容。普通对话为字符串;detail=true 或工作流命中交互节点时,可能为按字段名区分的对象数组(如 text / interactive / tool / file / reasoning)。当元素包含 interactive 字段时,形如 { interactive: { type, params } },只返回交互展示配置,不返回 entryNodeIds / memoryEdges / nodeOutputs 等内部运行态字段'
}), }),
reasoning_content: z.string().optional().meta({ description: '思考过程内容(仅推理模型有)' }) reasoning_content: z.string().optional().meta({ description: '思考过程内容(仅推理模型有)' })
}); });
......
...@@ -153,7 +153,7 @@ const detailTrueStreamFalseExample = { ...@@ -153,7 +153,7 @@ const detailTrueStreamFalseExample = {
] ]
}; };
// 交互节点-用户选择 (非流式响应,从 choices 中获取 type=interactive) // 交互节点-用户选择 (非流式响应,从 choices 中获取 interactive 字段)
const interactiveUserSelectResponseExample = { const interactiveUserSelectResponseExample = {
id: 'chatId', id: 'chatId',
model: '', model: '',
...@@ -164,7 +164,6 @@ const interactiveUserSelectResponseExample = { ...@@ -164,7 +164,6 @@ const interactiveUserSelectResponseExample = {
role: 'assistant', role: 'assistant',
content: [ content: [
{ {
type: 'interactive',
interactive: { interactive: {
type: 'userSelect', type: 'userSelect',
params: { params: {
...@@ -195,7 +194,6 @@ const interactiveUserInputResponseExample = { ...@@ -195,7 +194,6 @@ const interactiveUserInputResponseExample = {
role: 'assistant', role: 'assistant',
content: [ content: [
{ {
type: 'interactive',
interactive: { interactive: {
type: 'userInput', type: 'userInput',
params: { params: {
...@@ -356,7 +354,9 @@ export const ChatCompletionPath: OpenAPIPath = { ...@@ -356,7 +354,9 @@ export const ChatCompletionPath: OpenAPIPath = {
如果工作流中包含交互节点,需要设置 \`detail=true\`: 如果工作流中包含交互节点,需要设置 \`detail=true\`:
- \`stream=true\`:可从 \`event=interactive\` 数据中获取交互节点的配置。 - \`stream=true\`:可从 \`event=interactive\` 数据中获取交互节点的配置。
- \`stream=false\`:可从 \`choices\` 中获取 \`type=interactive\` 的元素。 - \`stream=false\`:可从 \`choices[].message.content\` 中获取包含 \`interactive\` 字段的元素。
返回给外部调用方的 \`interactive\` 是展示配置,只包含 \`type\` 和 \`params\`;\`entryNodeIds\` / \`memoryEdges\` / \`nodeOutputs\` / \`nodeResponseId\` 等内部运行态字段不会返回。若内部命中 children / loop / tool 包装交互,接口会返回最深层面向用户的交互节点。
接收到交互节点信息后,可以根据数据进行 UI 渲染并引导用户输入/选择,然后再次调用本接口继续工作流: 接收到交互节点信息后,可以根据数据进行 UI 渲染并引导用户输入/选择,然后再次调用本接口继续工作流:
...@@ -445,13 +445,13 @@ ${interactiveStreamExample} ...@@ -445,13 +445,13 @@ ${interactiveStreamExample}
interactiveUserSelect: { interactiveUserSelect: {
summary: '交互节点-用户选择 响应', summary: '交互节点-用户选择 响应',
description: description:
'工作流命中用户选择交互节点。从 choices[].message.content 中获取 type=interactive 的元素', '工作流命中用户选择交互节点。从 choices[].message.content 中获取包含 interactive 字段的元素',
value: interactiveUserSelectResponseExample value: interactiveUserSelectResponseExample
}, },
interactiveUserInput: { interactiveUserInput: {
summary: '交互节点-表单输入 响应', summary: '交互节点-表单输入 响应',
description: description:
'工作流命中表单输入交互节点。从 choices[].message.content 中获取 type=interactive 的元素', '工作流命中表单输入交互节点。从 choices[].message.content 中获取包含 interactive 字段的元素',
value: interactiveUserInputResponseExample value: interactiveUserInputResponseExample
} }
} }
...@@ -515,7 +515,9 @@ ${interactiveStreamExample} ...@@ -515,7 +515,9 @@ ${interactiveStreamExample}
如果工作流中包含交互节点,需要设置 \`detail=true\`: 如果工作流中包含交互节点,需要设置 \`detail=true\`:
- \`stream=true\`:可从 \`event=interactive\` 数据中获取交互节点的配置。 - \`stream=true\`:可从 \`event=interactive\` 数据中获取交互节点的配置。
- \`stream=false\`:可从 \`choices\` 中获取 \`type=interactive\` 的元素。 - \`stream=false\`:可从 \`choices[].message.content\` 中获取包含 \`interactive\` 字段的元素。
返回给外部调用方的 \`interactive\` 是展示配置,只包含 \`type\` 和 \`params\`;\`entryNodeIds\` / \`memoryEdges\` / \`nodeOutputs\` / \`nodeResponseId\` 等内部运行态字段不会返回。若内部命中 children / loop / tool 包装交互,接口会返回最深层面向用户的交互节点。
接收到交互节点信息后,可以根据数据进行 UI 渲染并引导用户输入/选择,然后再次调用本接口继续工作流: 接收到交互节点信息后,可以根据数据进行 UI 渲染并引导用户输入/选择,然后再次调用本接口继续工作流:
...@@ -604,13 +606,13 @@ ${interactiveStreamExample} ...@@ -604,13 +606,13 @@ ${interactiveStreamExample}
interactiveUserSelect: { interactiveUserSelect: {
summary: '交互节点-用户选择 响应', summary: '交互节点-用户选择 响应',
description: description:
'工作流命中用户选择交互节点。从 choices[].message.content 中获取 type=interactive 的元素', '工作流命中用户选择交互节点。从 choices[].message.content 中获取包含 interactive 字段的元素',
value: interactiveUserSelectResponseExample value: interactiveUserSelectResponseExample
}, },
interactiveUserInput: { interactiveUserInput: {
summary: '交互节点-表单输入 响应', summary: '交互节点-表单输入 响应',
description: description:
'工作流命中表单输入交互节点。从 choices[].message.content 中获取 type=interactive 的元素', '工作流命中表单输入交互节点。从 choices[].message.content 中获取包含 interactive 字段的元素',
value: interactiveUserInputResponseExample value: interactiveUserInputResponseExample
} }
} }
......
...@@ -3,15 +3,13 @@ import { ...@@ -3,15 +3,13 @@ import {
DispatchNodeResponseKeyEnum, DispatchNodeResponseKeyEnum,
SseResponseEventEnum SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants'; } from '@fastgpt/global/core/workflow/runtime/constants';
import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { UsageSourceEnum } from '@fastgpt/global/support/wallet/usage/constants'; import { UsageSourceEnum } from '@fastgpt/global/support/wallet/usage/constants';
import type { AIChatItemType, UserChatItemType } from '@fastgpt/global/core/chat/type'; import type { AIChatItemType, UserChatItemType } from '@fastgpt/global/core/chat/type';
import { GPTMessages2Chats } from '@fastgpt/global/core/chat/adapt'; import { GPTMessages2Chats } from '@fastgpt/global/core/chat/adapt';
import { concatHistories, removeEmptyUserInput } from '@fastgpt/global/core/chat/utils'; import { concatHistories, removeEmptyUserInput } from '@fastgpt/global/core/chat/utils';
import { WritePermissionVal } from '@fastgpt/global/support/permission/constant'; import { WritePermissionVal } from '@fastgpt/global/support/permission/constant';
import { import { getLastInteractiveValue } from '@fastgpt/global/core/workflow/runtime/utils';
getLastInteractiveValue,
textAdaptGptResponse
} from '@fastgpt/global/core/workflow/runtime/utils';
import { import {
ChatGenerateStatusEnum, ChatGenerateStatusEnum,
ChatRoleEnum, ChatRoleEnum,
...@@ -258,22 +256,11 @@ export async function handleSkillDebugChat( ...@@ -258,22 +256,11 @@ export async function handleSkillDebugChat(
logger.debug('Skill debug workflow completed', { skillId, chatId, durationSeconds }); logger.debug('Skill debug workflow completed', { skillId, chatId, durationSeconds });
computedFlowResponses.forEach((nodeResponse) => { computedFlowResponses.forEach((nodeResponse) => {
streamResponseContext?.responseWrite({ streamResponseContext?.responseWrite(workflowSseEvent.flowNodeResponse(nodeResponse));
event: SseResponseEventEnum.flowNodeResponse,
data: nodeResponse
});
});
streamResponseContext.responseWrite({
event: SseResponseEventEnum.workflowDuration,
data: {
durationSeconds
}
}); });
streamResponseContext.responseWrite(workflowSseEvent.workflowDuration(durationSeconds));
streamResponseContext.responseWrite({ streamResponseContext.responseWrite(workflowSseEvent.answerStop());
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({ text: null, finish_reason: 'stop' })
});
const aiResponse: AIChatItemType & { dataId?: string } = { const aiResponse: AIChatItemType & { dataId?: string } = {
dataId: finalResponseChatItemId, dataId: finalResponseChatItemId,
...@@ -319,10 +306,7 @@ export async function handleSkillDebugChat( ...@@ -319,10 +306,7 @@ export async function handleSkillDebugChat(
}); });
} }
streamResponseContext.responseWrite({ streamResponseContext.responseWrite(workflowSseEvent.done(SseResponseEventEnum.answer));
event: SseResponseEventEnum.answer,
data: '[DONE]'
});
await streamResponseContext.flushResume(); await streamResponseContext.flushResume();
} catch (err: any) { } catch (err: any) {
......
import type { UserChatItemType } from '@fastgpt/global/core/chat/type'; import type { UserChatItemType } from '@fastgpt/global/core/chat/type';
import { ChatCompletionRequestMessageRoleEnum } from '@fastgpt/global/core/ai/constants'; import { ChatCompletionRequestMessageRoleEnum } from '@fastgpt/global/core/ai/constants';
import { chatValue2RuntimePrompt } from '@fastgpt/global/core/chat/adapt'; import { chatValue2RuntimePrompt } from '@fastgpt/global/core/chat/adapt';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import type { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants';
import { delay, withTimeout } from '@fastgpt/global/common/system/utils'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import type { WorkflowTypedSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { withTimeout } from '@fastgpt/global/common/system/utils';
import { getLogger, LogCategories } from '../../common/logger'; import { getLogger, LogCategories } from '../../common/logger';
import { createLLMResponse } from '../ai/llm/request'; import { createLLMResponse } from '../ai/llm/request';
import { getDefaultChatTitleModel } from '../ai/model'; import { getDefaultChatTitleModel } from '../ai/model';
...@@ -271,10 +273,7 @@ export const createGeneratedChatTitleSender = ({ ...@@ -271,10 +273,7 @@ export const createGeneratedChatTitleSender = ({
titleGeneration?: Promise<GeneratedChatTitleResult | undefined>; titleGeneration?: Promise<GeneratedChatTitleResult | undefined>;
stream: boolean; stream: boolean;
detail: boolean; detail: boolean;
writeChatTitle?: (payload: { writeChatTitle?: (payload: WorkflowTypedSseEvent<SseResponseEventEnum.chatTitle>) => void;
event: SseResponseEventEnum.chatTitle;
data: { title: string };
}) => void;
}) => { }) => {
const titleResultPromise = titleGeneration?.catch(() => undefined); const titleResultPromise = titleGeneration?.catch(() => undefined);
let titleEventWritten = false; let titleEventWritten = false;
...@@ -304,12 +303,7 @@ export const createGeneratedChatTitleSender = ({ ...@@ -304,12 +303,7 @@ export const createGeneratedChatTitleSender = ({
const { title } = titleResult; const { title } = titleResult;
if (stream && detail && !titleEventWritten && !closed) { if (stream && detail && !titleEventWritten && !closed) {
writeChatTitle?.({ writeChatTitle?.(workflowSseEvent.chatTitle(title));
event: SseResponseEventEnum.chatTitle,
data: {
title
}
});
titleEventWritten = true; titleEventWritten = true;
} }
......
...@@ -4,12 +4,11 @@ import type { ModuleDispatchProps } from '@fastgpt/global/core/workflow/runtime/ ...@@ -4,12 +4,11 @@ import type { ModuleDispatchProps } from '@fastgpt/global/core/workflow/runtime/
import { type SelectAppItemType } from '@fastgpt/global/core/workflow/template/system/abandoned/runApp/type'; import { type SelectAppItemType } from '@fastgpt/global/core/workflow/template/system/abandoned/runApp/type';
import { runWorkflow } from '../index'; import { runWorkflow } from '../index';
import { ChatRoleEnum, ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants'; import { ChatRoleEnum, ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { import {
getWorkflowEntryNodeIds, getWorkflowEntryNodeIds,
storeEdges2RuntimeEdges, storeEdges2RuntimeEdges,
storeNodes2RuntimeNodes, storeNodes2RuntimeNodes
textAdaptGptResponse
} from '@fastgpt/global/core/workflow/runtime/utils'; } from '@fastgpt/global/core/workflow/runtime/utils';
import type { NodeInputKeyEnum, NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import type { NodeInputKeyEnum, NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
...@@ -52,12 +51,7 @@ export const dispatchAppRequest = async (props: Props): Promise<Response> => { ...@@ -52,12 +51,7 @@ export const dispatchAppRequest = async (props: Props): Promise<Response> => {
per: ReadPermissionVal per: ReadPermissionVal
}); });
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.fastAnswerDelta('\n'));
event: SseResponseEventEnum.fastAnswer,
data: textAdaptGptResponse({
text: '\n'
})
});
const chatHistories = getHistories(history, histories); const chatHistories = getHistories(history, histories);
const { files } = chatValue2RuntimePrompt(query); const { files } = chatValue2RuntimePrompt(query);
......
...@@ -2,8 +2,7 @@ import type { ...@@ -2,8 +2,7 @@ import type {
AIChatItemValueItemType, AIChatItemValueItemType,
ToolModuleResponseItemType ToolModuleResponseItemType
} from '@fastgpt/global/core/chat/type'; } from '@fastgpt/global/core/chat/type';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import type { AgentLoopEvent } from '../../../../../ai/llm/agentLoop'; import type { AgentLoopEvent } from '../../../../../ai/llm/agentLoop';
import type { WorkflowResponseType } from '../../../type'; import type { WorkflowResponseType } from '../../../type';
import type { GetSubAppInfoFnType } from '../type'; import type { GetSubAppInfoFnType } from '../type';
...@@ -270,16 +269,7 @@ export const createWorkflowAgentLoopEventMapper = ({ ...@@ -270,16 +269,7 @@ export const createWorkflowAgentLoopEventMapper = ({
params: `${tool.params || ''}${argsDelta}` params: `${tool.params || ''}${argsDelta}`
})); }));
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.toolParams({ id: callId, params: argsDelta }));
id: callId,
event: SseResponseEventEnum.toolParams,
data: {
tool: {
id: callId,
params: argsDelta
}
}
});
}; };
const applyToolResponse = ({ callId, response }: { callId: string; response: string }) => { const applyToolResponse = ({ callId, response }: { callId: string; response: string }) => {
...@@ -298,16 +288,7 @@ export const createWorkflowAgentLoopEventMapper = ({ ...@@ -298,16 +288,7 @@ export const createWorkflowAgentLoopEventMapper = ({
response: `${tool.response || ''}${response}` response: `${tool.response || ''}${response}`
})); }));
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.toolResponse({ id: callId, response }));
id: callId,
event: SseResponseEventEnum.toolResponse,
data: {
tool: {
id: callId,
response
}
}
});
}; };
/** /**
...@@ -316,33 +297,17 @@ export const createWorkflowAgentLoopEventMapper = ({ ...@@ -316,33 +297,17 @@ export const createWorkflowAgentLoopEventMapper = ({
const emitEvent = (event: AgentLoopEvent) => { const emitEvent = (event: AgentLoopEvent) => {
switch (event.type) { switch (event.type) {
case 'answer_delta': { case 'answer_delta': {
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.answerDelta(event.text));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
text: event.text
})
});
return; return;
} }
case 'reasoning_delta': { case 'reasoning_delta': {
if (!showReasoning) return; if (!showReasoning) return;
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.reasoningDelta(event.text));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
reasoning_content: event.text
})
});
return; return;
} }
case 'llm_request_start': { case 'llm_request_start': {
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.flowNodeStatus(event.modelName));
event: SseResponseEventEnum.flowNodeStatus,
data: {
status: 'running',
name: event.modelName
}
});
return; return;
} }
case 'llm_request_end': { case 'llm_request_end': {
...@@ -412,13 +377,7 @@ export const createWorkflowAgentLoopEventMapper = ({ ...@@ -412,13 +377,7 @@ export const createWorkflowAgentLoopEventMapper = ({
}; };
upsertToolResponse(tool); upsertToolResponse(tool);
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.toolCall(tool));
id: event.call.id,
event: SseResponseEventEnum.toolCall,
data: {
tool
}
});
return; return;
} }
case 'tool_params': { case 'tool_params': {
...@@ -449,15 +408,9 @@ export const createWorkflowAgentLoopEventMapper = ({ ...@@ -449,15 +408,9 @@ export const createWorkflowAgentLoopEventMapper = ({
return; return;
} }
case 'plan_status': { case 'plan_status': {
workflowStreamResponse?.({ workflowStreamResponse?.(
id: AGENT_PLAN_STREAM_RESPONSE_ID, workflowSseEvent.planStatus({ status: event.status }, AGENT_PLAN_STREAM_RESPONSE_ID)
event: SseResponseEventEnum.planStatus, );
data: {
planStatus: {
status: event.status
}
}
});
return; return;
} }
case 'plan_update': { case 'plan_update': {
...@@ -476,13 +429,7 @@ export const createWorkflowAgentLoopEventMapper = ({ ...@@ -476,13 +429,7 @@ export const createWorkflowAgentLoopEventMapper = ({
assistantResponses.push(nextPlanValue); assistantResponses.push(nextPlanValue);
} }
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.plan(event.plan, AGENT_PLAN_STREAM_RESPONSE_ID));
id: AGENT_PLAN_STREAM_RESPONSE_ID,
event: SseResponseEventEnum.plan,
data: {
plan: event.plan
}
});
return; return;
} }
} }
......
...@@ -6,8 +6,7 @@ import type { ...@@ -6,8 +6,7 @@ import type {
ChatCompletionMessageToolCall, ChatCompletionMessageToolCall,
CompletionFinishReason CompletionFinishReason
} from '@fastgpt/global/core/ai/llm/type'; } from '@fastgpt/global/core/ai/llm/type';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import type { ChatNodeUsageType } from '@fastgpt/global/support/wallet/bill/type'; import type { ChatNodeUsageType } from '@fastgpt/global/support/wallet/bill/type';
import type { AgentEvent, AgentMessage } from '@mariozechner/pi-agent-core'; import type { AgentEvent, AgentMessage } from '@mariozechner/pi-agent-core';
import type { AssistantMessage, Model, StopReason, ToolCall } from '@mariozechner/pi-ai'; import type { AssistantMessage, Model, StopReason, ToolCall } from '@mariozechner/pi-ai';
...@@ -419,13 +418,7 @@ export const createPiAgentWorkflowRuntime = ({ ...@@ -419,13 +418,7 @@ export const createPiAgentWorkflowRuntime = ({
}; };
pendingRequests.push(request); pendingRequests.push(request);
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.flowNodeStatus(request.modelName));
event: SseResponseEventEnum.flowNodeStatus,
data: {
status: 'running',
name: request.modelName
}
});
return undefined; return undefined;
}, },
...@@ -434,19 +427,13 @@ export const createPiAgentWorkflowRuntime = ({ ...@@ -434,19 +427,13 @@ export const createPiAgentWorkflowRuntime = ({
const assistantEvent = event.assistantMessageEvent; const assistantEvent = event.assistantMessageEvent;
if (assistantEvent.type === 'text_delta') { if (assistantEvent.type === 'text_delta') {
answerText += assistantEvent.delta; answerText += assistantEvent.delta;
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.answerDelta(assistantEvent.delta));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({ text: assistantEvent.delta })
});
return; return;
} }
if (assistantEvent.type === 'thinking_delta') { if (assistantEvent.type === 'thinking_delta') {
reasoningText += assistantEvent.delta; reasoningText += assistantEvent.delta;
if (showReasoning) { if (showReasoning) {
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.reasoningDelta(assistantEvent.delta));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({ reasoning_content: assistantEvent.delta })
});
} }
return; return;
} }
......
...@@ -3,7 +3,7 @@ import type { ...@@ -3,7 +3,7 @@ import type {
ChatHistoryItemResType, ChatHistoryItemResType,
ToolModuleResponseItemType ToolModuleResponseItemType
} from '@fastgpt/global/core/chat/type'; } from '@fastgpt/global/core/chat/type';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant'; import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import type { DispatchAgentModuleProps } from '..'; import type { DispatchAgentModuleProps } from '..';
import { getExecuteTool, type ToolDispatchContext } from '../utils'; import { getExecuteTool, type ToolDispatchContext } from '../utils';
...@@ -171,28 +171,13 @@ export const createPiAgentToolEventHandler = ({ ...@@ -171,28 +171,13 @@ export const createPiAgentToolEventHandler = ({
}; };
upsertAssistantTool(assistantResponses, assistantTool); upsertAssistantTool(assistantResponses, assistantTool);
ctx.streamResponseFn?.({ ctx.streamResponseFn?.(workflowSseEvent.toolCall(assistantTool));
id: callId,
event: SseResponseEventEnum.toolCall,
data: {
tool: assistantTool
}
});
} }
const latestTool = findAssistantTool(assistantResponses, callId); const latestTool = findAssistantTool(assistantResponses, callId);
if (argStr && !latestTool?.params) { if (argStr && !latestTool?.params) {
appendAssistantToolParams(assistantResponses, callId, argStr); appendAssistantToolParams(assistantResponses, callId, argStr);
ctx.streamResponseFn?.({ ctx.streamResponseFn?.(workflowSseEvent.toolParams({ id: callId, params: argStr }));
id: callId,
event: SseResponseEventEnum.toolParams,
data: {
tool: {
id: callId,
params: argStr
}
}
});
} }
}; };
...@@ -257,16 +242,12 @@ export const createPiAgentToolEventHandler = ({ ...@@ -257,16 +242,12 @@ export const createPiAgentToolEventHandler = ({
const currentTool = findAssistantTool(assistantResponses, event.toolCallId); const currentTool = findAssistantTool(assistantResponses, event.toolCallId);
if (response && !currentTool?.response) { if (response && !currentTool?.response) {
appendAssistantToolResponse(assistantResponses, event.toolCallId, response); appendAssistantToolResponse(assistantResponses, event.toolCallId, response);
ctx.streamResponseFn?.({ ctx.streamResponseFn?.(
id: event.toolCallId, workflowSseEvent.toolResponse({
event: SseResponseEventEnum.toolResponse, id: event.toolCallId,
data: { response
tool: { })
id: event.toolCallId, );
response
}
}
});
} }
if (event.isError) { if (event.isError) {
...@@ -336,16 +317,7 @@ export async function buildAgentTools({ ...@@ -336,16 +317,7 @@ export async function buildAgentTools({
if (usages.length > 0) usagePush(usages); if (usages.length > 0) usagePush(usages);
appendAssistantToolResponse(assistantResponses, callId, response); appendAssistantToolResponse(assistantResponses, callId, response);
ctx.streamResponseFn?.({ ctx.streamResponseFn?.(workflowSseEvent.toolResponse({ id: callId, response }));
id: callId,
event: SseResponseEventEnum.toolResponse,
data: {
tool: {
id: callId,
response
}
}
});
return { content: [{ type: 'text' as const, text: response }], details: {} }; return { content: [{ type: 'text' as const, text: response }], details: {} };
}; };
...@@ -364,27 +336,12 @@ export async function buildAgentTools({ ...@@ -364,27 +336,12 @@ export async function buildAgentTools({
}; };
upsertAssistantTool(assistantResponses, assistantTool); upsertAssistantTool(assistantResponses, assistantTool);
ctx.streamResponseFn?.({ ctx.streamResponseFn?.(workflowSseEvent.toolCall(assistantTool));
id: callId,
event: SseResponseEventEnum.toolCall,
data: {
tool: assistantTool
}
});
} }
if (argStr && !findAssistantTool(assistantResponses, callId)?.params) { if (argStr && !findAssistantTool(assistantResponses, callId)?.params) {
appendAssistantToolParams(assistantResponses, callId, argStr); appendAssistantToolParams(assistantResponses, callId, argStr);
ctx.streamResponseFn?.({ ctx.streamResponseFn?.(workflowSseEvent.toolParams({ id: callId, params: argStr }));
id: callId,
event: SseResponseEventEnum.toolParams,
data: {
tool: {
id: callId,
params: argStr
}
}
});
} }
return execute(callId, args, argStr); return execute(callId, args, argStr);
......
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import type { WorkflowResponseType } from '../../../../type'; import type { WorkflowResponseType } from '../../../../type';
import { getRunningSandboxId } from '../../../../../../ai/sandbox/interface/runtime'; import { getRunningSandboxId } from '../../../../../../ai/sandbox/interface/runtime';
import type { ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants'; import type { ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants';
...@@ -28,11 +28,10 @@ export const streamAgentSandboxInitStatus = ({ ...@@ -28,11 +28,10 @@ export const streamAgentSandboxInitStatus = ({
chatId chatId
}); });
workflowStreamResponse?.({ workflowStreamResponse?.(
event: SseResponseEventEnum.sandboxStatus, workflowSseEvent.sandboxStatus({
data: {
sandboxId: effectiveSandboxId, sandboxId: effectiveSandboxId,
phase: 'lazyInit' phase: 'lazyInit'
} })
}); );
}; };
...@@ -11,8 +11,7 @@ import type { ...@@ -11,8 +11,7 @@ import type {
ChatDispatchProps, ChatDispatchProps,
RuntimeNodeItemType RuntimeNodeItemType
} from '@fastgpt/global/core/workflow/runtime/type'; } from '@fastgpt/global/core/workflow/runtime/type';
import type { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { pushTrack } from '../../../../../../../common/middle/tracks/utils'; import { pushTrack } from '../../../../../../../common/middle/tracks/utils';
...@@ -20,7 +19,7 @@ import { getErrText } from '@fastgpt/global/common/error/utils'; ...@@ -20,7 +19,7 @@ import { getErrText } from '@fastgpt/global/common/error/utils';
import { getAppVersionById } from '../../../../../../app/version/controller'; import { getAppVersionById } from '../../../../../../app/version/controller';
import { assertMCPUrlNotInternal, MCPClient } from '../../../../../../app/mcp'; import { assertMCPUrlNotInternal, MCPClient } from '../../../../../../app/mcp';
import { runHTTPTool } from '../../../../../../app/http'; import { runHTTPTool } from '../../../../../../app/http';
import { parseToolId } from '../../../../child/runTool'; import { isPluginAnswerType, parseToolId } from '../../../../child/runTool';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant'; import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import type { RequireOnlyOne } from '@fastgpt/global/common/type/utils'; import type { RequireOnlyOne } from '@fastgpt/global/common/type/utils';
import { pluginClient } from '../../../../../../../thirdProvider/fastgptPlugin'; import { pluginClient } from '../../../../../../../thirdProvider/fastgptPlugin';
...@@ -170,15 +169,14 @@ export const dispatchTool = async ({ ...@@ -170,15 +169,14 @@ export const dispatchTool = async ({
invokeToken: invokeToken || '' invokeToken: invokeToken || ''
}, },
onMessage: ({ type, content }) => { onMessage: ({ type, content }) => {
if (workflowStreamResponse && content) { if (!workflowStreamResponse || !content || !isPluginAnswerType(type)) return;
answerText += content;
workflowStreamResponse({ answerText += content;
event: type as unknown as SseResponseEventEnum, workflowStreamResponse(
data: textAdaptGptResponse({ type === 'fastAnswer'
text: content ? workflowSseEvent.fastAnswerDelta(content)
}) : workflowSseEvent.answerDelta(content)
}); );
}
} }
}); });
......
...@@ -2,11 +2,8 @@ import { i18nT } from '@fastgpt/global/common/i18n/utils'; ...@@ -2,11 +2,8 @@ import { i18nT } from '@fastgpt/global/common/i18n/utils';
import { getQuoteTemplate } from '@fastgpt/global/core/ai/prompt/AIChat'; import { getQuoteTemplate } from '@fastgpt/global/core/ai/prompt/AIChat';
import { GPTMessages2Chats } from '@fastgpt/global/core/chat/adapt'; import { GPTMessages2Chats } from '@fastgpt/global/core/chat/adapt';
import { getHistoryPreview } from '@fastgpt/global/core/chat/utils'; import { getHistoryPreview } from '@fastgpt/global/core/chat/utils';
import { import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
DispatchNodeResponseKeyEnum, import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import { getLLMModel } from '../../../../ai/model'; import { getLLMModel } from '../../../../ai/model';
import { createLLMResponse } from '../../../../ai/llm/request'; import { createLLMResponse } from '../../../../ai/llm/request';
import { computedMaxToken } from '../../../../ai/utils'; import { computedMaxToken } from '../../../../ai/utils';
...@@ -173,21 +170,11 @@ export const dispatchChatCompletion = async (props: ChatProps): Promise<ChatResp ...@@ -173,21 +170,11 @@ export const dispatchChatCompletion = async (props: ChatProps): Promise<ChatResp
isAborted: checkIsStopping, isAborted: checkIsStopping,
onReasoning({ text }) { onReasoning({ text }) {
if (!aiChatReasoning) return; if (!aiChatReasoning) return;
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.reasoningDelta(text));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
reasoning_content: text
})
});
}, },
onStreaming({ text }) { onStreaming({ text }) {
if (!isResponseAnswerText) return; if (!isResponseAnswerText) return;
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.answerDelta(text));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
text
})
});
} }
}); });
......
import type { ChatCompletionMessageToolCall } from '@fastgpt/global/core/ai/llm/type'; import type { ChatCompletionMessageToolCall } from '@fastgpt/global/core/ai/llm/type';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import { sliceStrStartEnd } from '@fastgpt/global/common/string/tools'; import { sliceStrStartEnd } from '@fastgpt/global/common/string/tools';
import type { WorkflowResponseType } from '../../../type'; import type { WorkflowResponseType } from '../../../type';
import type { ToolInfo } from './useToolCatalog'; import type { ToolInfo } from './useToolCatalog';
...@@ -18,22 +17,12 @@ export const useToolStreamResponse = ({ ...@@ -18,22 +17,12 @@ export const useToolStreamResponse = ({
}) => { }) => {
const streamReasoning = (text: string) => { const streamReasoning = (text: string) => {
if (!aiChatReasoning) return; if (!aiChatReasoning) return;
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.reasoningDelta(text));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
reasoning_content: text
})
});
}; };
const streamAnswer = (text: string) => { const streamAnswer = (text: string) => {
if (!isResponseAnswerText) return; if (!isResponseAnswerText) return;
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.answerDelta(text));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
text
})
});
}; };
const streamToolCall = (call: ChatCompletionMessageToolCall) => { const streamToolCall = (call: ChatCompletionMessageToolCall) => {
...@@ -41,19 +30,15 @@ export const useToolStreamResponse = ({ ...@@ -41,19 +30,15 @@ export const useToolStreamResponse = ({
const toolNode = getToolInfo(call.function.name); const toolNode = getToolInfo(call.function.name);
if (!toolNode) return; if (!toolNode) return;
workflowStreamResponse?.({ workflowStreamResponse?.(
id: call.id, workflowSseEvent.toolCall({
event: SseResponseEventEnum.toolCall, id: call.id,
data: { toolName: toolNode.name,
tool: { toolAvatar: toolNode.avatar ?? '',
id: call.id, functionName: call.function.name,
toolName: toolNode.name, params: call.function.arguments ?? ''
toolAvatar: toolNode.avatar, })
functionName: call.function.name, );
params: call.function.arguments ?? ''
}
}
});
}; };
const streamToolParams = ({ const streamToolParams = ({
...@@ -64,18 +49,14 @@ export const useToolStreamResponse = ({ ...@@ -64,18 +49,14 @@ export const useToolStreamResponse = ({
argsDelta: string; argsDelta: string;
}) => { }) => {
if (!isResponseAnswerText) return; if (!isResponseAnswerText) return;
workflowStreamResponse?.({ workflowStreamResponse?.(
id: call.id, workflowSseEvent.toolParams({
event: SseResponseEventEnum.toolParams, id: call.id,
data: { toolName: '',
tool: { toolAvatar: '',
id: call.id, params: argsDelta
toolName: '', })
toolAvatar: '', );
params: argsDelta
}
}
});
}; };
const streamToolResponse = ({ const streamToolResponse = ({
...@@ -90,19 +71,15 @@ export const useToolStreamResponse = ({ ...@@ -90,19 +71,15 @@ export const useToolStreamResponse = ({
/** /**
* SSE 只给聊天气泡做轻量预览;完整工具响应与压缩 child 保存在 nodeResponse。 * SSE 只给聊天气泡做轻量预览;完整工具响应与压缩 child 保存在 nodeResponse。
*/ */
workflowStreamResponse?.({ workflowStreamResponse?.(
id: toolCallId, workflowSseEvent.toolResponse({
event: SseResponseEventEnum.toolResponse, id: toolCallId,
data: { toolName: '',
tool: { toolAvatar: '',
id: toolCallId, params: '',
toolName: '', response: sliceStrStartEnd(response || '', 5000, 5000)
toolAvatar: '', })
params: '', );
response: sliceStrStartEnd(response || '', 5000, 5000)
}
}
});
}; };
return { return {
......
...@@ -2,13 +2,12 @@ import type { ChatItemMiniType } from '@fastgpt/global/core/chat/type'; ...@@ -2,13 +2,12 @@ import type { ChatItemMiniType } from '@fastgpt/global/core/chat/type';
import type { ModuleDispatchProps } from '@fastgpt/global/core/workflow/runtime/type'; import type { ModuleDispatchProps } from '@fastgpt/global/core/workflow/runtime/type';
import { runWorkflow } from '../index'; import { runWorkflow } from '../index';
import { ChatRoleEnum, ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants'; import { ChatRoleEnum, ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { import {
getWorkflowEntryNodeIds, getWorkflowEntryNodeIds,
storeEdges2RuntimeEdges, storeEdges2RuntimeEdges,
rewriteNodeOutputByHistories, rewriteNodeOutputByHistories,
storeNodes2RuntimeNodes, storeNodes2RuntimeNodes
textAdaptGptResponse
} from '@fastgpt/global/core/workflow/runtime/utils'; } from '@fastgpt/global/core/workflow/runtime/utils';
import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
...@@ -94,12 +93,7 @@ export const dispatchRunAppNode = async (props: Props): Promise<Response> => { ...@@ -94,12 +93,7 @@ export const dispatchRunAppNode = async (props: Props): Promise<Response> => {
const childStreamResponse = system_forbid_stream ? false : props.stream; const childStreamResponse = system_forbid_stream ? false : props.stream;
// Auto line // Auto line
if (childStreamResponse) { if (childStreamResponse) {
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.answerDelta('\n'));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
text: '\n'
})
});
} }
const chatHistories = getHistories(history, histories); const chatHistories = getHistories(history, histories);
......
import { getErrText } from '@fastgpt/global/common/error/utils'; import { getErrText } from '@fastgpt/global/common/error/utils';
import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import type { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants';
import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { import {
type DispatchNodeResultType, type DispatchNodeResultType,
type ModuleDispatchProps type ModuleDispatchProps
...@@ -13,7 +13,6 @@ import type { McpToolDataType } from '@fastgpt/global/core/app/tool/mcpTool/type ...@@ -13,7 +13,6 @@ import type { McpToolDataType } from '@fastgpt/global/core/app/tool/mcpTool/type
import type { HttpToolConfigType } from '@fastgpt/global/core/app/tool/httpTool/type'; import type { HttpToolConfigType } from '@fastgpt/global/core/app/tool/httpTool/type';
import { SystemToolSecretInputTypeEnum } from '@fastgpt/global/core/app/tool/systemTool/constants'; import { SystemToolSecretInputTypeEnum } from '@fastgpt/global/core/app/tool/systemTool/constants';
import type { StoreSecretValueType } from '@fastgpt/global/common/secret/type'; import type { StoreSecretValueType } from '@fastgpt/global/common/secret/type';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import { pushTrack } from '../../../../common/middle/tracks/utils'; import { pushTrack } from '../../../../common/middle/tracks/utils';
import { getNodeErrResponse } from '../utils'; import { getNodeErrResponse } from '../utils';
import { getAppVersionById } from '../../../../core/app/version/controller'; import { getAppVersionById } from '../../../../core/app/version/controller';
...@@ -37,6 +36,9 @@ type SystemInputConfigType = { ...@@ -37,6 +36,9 @@ type SystemInputConfigType = {
value: StoreSecretValueType; value: StoreSecretValueType;
}; };
export const isPluginAnswerType = (type: string): type is 'answer' | 'fastAnswer' =>
type === 'answer' || type === 'fastAnswer';
type RunToolProps = ModuleDispatchProps<{ type RunToolProps = ModuleDispatchProps<{
[NodeInputKeyEnum.toolData]?: McpToolDataType; [NodeInputKeyEnum.toolData]?: McpToolDataType;
[NodeInputKeyEnum.systemInputConfig]?: SystemInputConfigType; [NodeInputKeyEnum.systemInputConfig]?: SystemInputConfigType;
...@@ -153,15 +155,14 @@ export const dispatchRunTool = async (props: RunToolProps): Promise<RunToolRespo ...@@ -153,15 +155,14 @@ export const dispatchRunTool = async (props: RunToolProps): Promise<RunToolRespo
time: cTime time: cTime
}, },
onMessage: ({ type, content }) => { onMessage: ({ type, content }) => {
if (workflowStreamResponse && content) { if (!workflowStreamResponse || !content || !isPluginAnswerType(type)) return;
answerText += content;
workflowStreamResponse({ answerText += content;
event: type as unknown as SseResponseEventEnum, workflowStreamResponse(
data: textAdaptGptResponse({ type === 'fastAnswer'
text: content ? workflowSseEvent.fastAnswerDelta(content)
}) : workflowSseEvent.answerDelta(content)
}); );
}
} }
}); });
......
...@@ -12,10 +12,8 @@ import type { ...@@ -12,10 +12,8 @@ import type {
} from '@fastgpt/global/core/workflow/runtime/type'; } from '@fastgpt/global/core/workflow/runtime/type';
import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant'; import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import { import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
DispatchNodeResponseKeyEnum, import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import { normalizeAIChatValue } from '@fastgpt/global/core/chat/adapt'; import { normalizeAIChatValue } from '@fastgpt/global/core/chat/adapt';
import type { import type {
ChatDispatchProps, ChatDispatchProps,
...@@ -804,13 +802,7 @@ export class WorkflowQueue { ...@@ -804,13 +802,7 @@ export class WorkflowQueue {
: getNanoid(); : getNanoid();
// push run status messages // push run status messages
if (node.showStatus && !this.data.isToolCall) { if (node.showStatus && !this.data.isToolCall) {
this.data.workflowStreamResponse?.({ this.data.workflowStreamResponse?.(workflowSseEvent.flowNodeStatus(node.name));
event: SseResponseEventEnum.flowNodeStatus,
data: {
status: 'running',
name: node.name
}
});
} }
const startTime = Date.now(); const startTime = Date.now();
// get node running params // get node running params
...@@ -993,10 +985,7 @@ export class WorkflowQueue { ...@@ -993,10 +985,7 @@ export class WorkflowQueue {
}); });
filteredResponses.forEach((item) => { filteredResponses.forEach((item) => {
this.data.workflowStreamResponse?.({ this.data.workflowStreamResponse?.(workflowSseEvent.flowNodeResponse(item));
event: SseResponseEventEnum.flowNodeResponse,
data: item
});
}); });
} }
...@@ -1466,10 +1455,7 @@ export class WorkflowQueue { ...@@ -1466,10 +1455,7 @@ export class WorkflowQueue {
// Tool call, not need interactive response // Tool call, not need interactive response
if (!this.data.isToolCall && this.isRootRuntime) { if (!this.data.isToolCall && this.isRootRuntime) {
this.data.workflowStreamResponse?.({ this.data.workflowStreamResponse?.(workflowSseEvent.interactive(interactiveResult));
event: SseResponseEventEnum.interactive,
data: { interactive: interactiveResult }
});
} }
return { return {
...@@ -1621,10 +1607,7 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR ...@@ -1621,10 +1607,7 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR
workflowSpan.setStatus({ code: SpanStatusCode.OK }); workflowSpan.setStatus({ code: SpanStatusCode.OK });
if (isRootRuntime) { if (isRootRuntime) {
data.workflowStreamResponse?.({ data.workflowStreamResponse?.(workflowSseEvent.workflowDuration(durationSeconds));
event: SseResponseEventEnum.workflowDuration,
data: { durationSeconds }
});
} }
return { return {
......
import { cloneDeep } from 'lodash'; import { cloneDeep } from 'lodash';
import { getErrText } from '@fastgpt/global/common/error/utils'; import { getErrText } from '@fastgpt/global/common/error/utils';
import { NodeInputKeyEnum, NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import { NodeInputKeyEnum, NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
DispatchNodeResponseKeyEnum, import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant'; import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import type { import type {
DispatchNodeResultType, DispatchNodeResultType,
...@@ -258,10 +256,7 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => { ...@@ -258,10 +256,7 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
); );
if (props.apiVersion === 'v2') { if (props.apiVersion === 'v2') {
recordedWrappers.forEach((item) => { recordedWrappers.forEach((item) => {
props.workflowStreamResponse?.({ props.workflowStreamResponse?.(workflowSseEvent.flowNodeResponse(item));
event: SseResponseEventEnum.flowNodeResponse,
data: item
});
}); });
} }
} }
......
import { batchRun } from '@fastgpt/global/common/system/utils'; import { batchRun } from '@fastgpt/global/common/system/utils';
import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
DispatchNodeResponseKeyEnum, import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import { import {
type DispatchNodeResultType, type DispatchNodeResultType,
type ModuleDispatchProps type ModuleDispatchProps
...@@ -199,10 +197,7 @@ export const dispatchParallelRun = async (props: Props): Promise<Response> => { ...@@ -199,10 +197,7 @@ export const dispatchParallelRun = async (props: Props): Promise<Response> => {
); );
if (props.apiVersion === 'v2') { if (props.apiVersion === 'v2') {
recordedWrappers.forEach((item) => { recordedWrappers.forEach((item) => {
props.workflowStreamResponse?.({ props.workflowStreamResponse?.(workflowSseEvent.flowNodeResponse(item));
event: SseResponseEventEnum.flowNodeResponse,
data: item
});
}); });
} }
} }
......
import { import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
DispatchNodeResponseKeyEnum, import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import type { ModuleDispatchProps } from '@fastgpt/global/core/workflow/runtime/type'; import type { ModuleDispatchProps } from '@fastgpt/global/core/workflow/runtime/type';
import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import { NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { type DispatchNodeResultType } from '@fastgpt/global/core/workflow/runtime/type'; import { type DispatchNodeResultType } from '@fastgpt/global/core/workflow/runtime/type';
...@@ -22,12 +19,7 @@ export const dispatchAnswer = (props: Record<string, any>): AnswerResponse => { ...@@ -22,12 +19,7 @@ export const dispatchAnswer = (props: Record<string, any>): AnswerResponse => {
const formatText = typeof text === 'string' ? text : JSON.stringify(text, null, 2); const formatText = typeof text === 'string' ? text : JSON.stringify(text, null, 2);
const responseText = `\n${formatText}`; const responseText = `\n${formatText}`;
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.fastAnswerDelta(responseText));
event: SseResponseEventEnum.fastAnswer,
data: textAdaptGptResponse({
text: responseText
})
});
return { return {
data: { data: {
......
import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants'; import type { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants';
import { VARIABLE_NODE_ID, WorkflowIOValueTypeEnum } from '@fastgpt/global/core/workflow/constants'; import { VARIABLE_NODE_ID, WorkflowIOValueTypeEnum } from '@fastgpt/global/core/workflow/constants';
import { FlowNodeInputTypeEnum } from '@fastgpt/global/core/workflow/node/constant'; import { FlowNodeInputTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import { import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
DispatchNodeResponseKeyEnum, import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import { type DispatchNodeResultType } from '@fastgpt/global/core/workflow/runtime/type'; import { type DispatchNodeResultType } from '@fastgpt/global/core/workflow/runtime/type';
import { getReferenceVariableValue } from '@fastgpt/global/core/workflow/runtime/utils'; import { getReferenceVariableValue } from '@fastgpt/global/core/workflow/runtime/utils';
import { type TUpdateListItem } from '@fastgpt/global/core/workflow/template/system/variableUpdate/type'; import { type TUpdateListItem } from '@fastgpt/global/core/workflow/template/system/variableUpdate/type';
...@@ -172,10 +170,7 @@ export const dispatchUpdateVariable = async (props: Props): Promise<Response> => ...@@ -172,10 +170,7 @@ export const dispatchUpdateVariable = async (props: Props): Promise<Response> =>
} }
if (!runningAppInfo.isChildApp) { if (!runningAppInfo.isChildApp) {
workflowStreamResponse?.({ workflowStreamResponse?.(workflowSseEvent.updateVariables(variableState.toStoreRecord()));
event: SseResponseEventEnum.updateVariables,
data: variableState.toStoreRecord()
});
} }
return { return {
......
...@@ -3,19 +3,21 @@ import type { ...@@ -3,19 +3,21 @@ import type {
ChatHistoryItemResType, ChatHistoryItemResType,
ToolRunResponseItemType ToolRunResponseItemType
} from '@fastgpt/global/core/chat/type'; } from '@fastgpt/global/core/chat/type';
import type { import type { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
DispatchNodeResponseKeyEnum,
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import type { RuntimeNodeItemType } from '@fastgpt/global/core/workflow/runtime/type'; import type { RuntimeNodeItemType } from '@fastgpt/global/core/workflow/runtime/type';
import type { import type {
WorkflowResponseItemType,
WorkflowResponseType
} from '@fastgpt/global/core/workflow/runtime/sse';
import type {
InteractiveNodeResponseType, InteractiveNodeResponseType,
WorkflowInteractiveResponseType WorkflowInteractiveResponseType
} from '@fastgpt/global/core/workflow/template/system/interactive/type'; } from '@fastgpt/global/core/workflow/template/system/interactive/type';
import { type RuntimeEdgeItemType } from '@fastgpt/global/core/workflow/type/edge'; import { type RuntimeEdgeItemType } from '@fastgpt/global/core/workflow/type/edge';
import type { ChatNodeUsageType } from '@fastgpt/global/support/wallet/bill/type'; import type { ChatNodeUsageType } from '@fastgpt/global/support/wallet/bill/type';
import type { NodeResponseWriteSummary } from '../../chat/nodeResponseStorage'; import type { NodeResponseWriteSummary } from '../../chat/nodeResponseStorage';
import z from 'zod';
export type { WorkflowResponseItemType, WorkflowResponseType };
/** /**
* workflow 内部运行期使用的 nodeResponse 摘要。 * workflow 内部运行期使用的 nodeResponse 摘要。
...@@ -90,16 +92,3 @@ export type DispatchFlowResponse = { ...@@ -90,16 +92,3 @@ export type DispatchFlowResponse = {
runtimeNodeResponseSummary: RuntimeNodeResponseSummary; runtimeNodeResponseSummary: RuntimeNodeResponseSummary;
durationSeconds: number; durationSeconds: number;
}; };
const WorkflowResponseItemSchema = z.object({
id: z.string().optional(),
event: z.custom<SseResponseEventEnum>().optional(),
data: z.union([z.string(), z.looseObject({})])
});
export type WorkflowResponseItemType = z.infer<typeof WorkflowResponseItemSchema>;
export const WorkflowResponseFnSchema = z.function({
input: z.tuple([WorkflowResponseItemSchema]),
output: z.void()
});
export type WorkflowResponseType = z.infer<typeof WorkflowResponseFnSchema>;
...@@ -3,8 +3,7 @@ import { getSseErrorResponse } from '../../../common/response'; ...@@ -3,8 +3,7 @@ import { getSseErrorResponse } from '../../../common/response';
import { clearCookie } from '../../../support/permission/auth/common'; import { clearCookie } from '../../../support/permission/auth/common';
import { STREAM_RESUME_REQUEST_HEADER } from '@fastgpt/global/core/chat/constants'; import { STREAM_RESUME_REQUEST_HEADER } from '@fastgpt/global/core/chat/constants';
import type { ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants'; import type { ChatSourceTypeEnum } from '@fastgpt/global/core/chat/constants';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { textAdaptGptResponse } from '@fastgpt/global/core/workflow/runtime/utils';
import { createSseStreamContext } from '../../../common/response/sse'; import { createSseStreamContext } from '../../../common/response/sse';
import { getStreamResumeMirror } from '../../chat/resume'; import { getStreamResumeMirror } from '../../chat/resume';
import { getWorkflowResponseWrite } from '../dispatch/utils'; import { getWorkflowResponseWrite } from '../dispatch/utils';
...@@ -110,12 +109,7 @@ export const initWorkflowSseResponse = ({ ...@@ -110,12 +109,7 @@ export const initWorkflowSseResponse = ({
// 10s 发送一次空 answer,沿用统一 SSE writer,避免浏览器或代理认为长连接已断开。 // 10s 发送一次空 answer,沿用统一 SSE writer,避免浏览器或代理认为长连接已断开。
heartbeat: { heartbeat: {
write: () => { write: () => {
responseWrite?.({ responseWrite?.(workflowSseEvent.answerDelta(''));
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
text: ''
})
});
} }
} }
}); });
......
...@@ -23,6 +23,7 @@ import { ...@@ -23,6 +23,7 @@ import {
SseResponseEventEnum, SseResponseEventEnum,
DispatchNodeResponseKeyEnum DispatchNodeResponseKeyEnum
} from '@fastgpt/global/core/workflow/runtime/constants'; } from '@fastgpt/global/core/workflow/runtime/constants';
import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant'; import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import type { RuntimeEdgeItemType } from '@fastgpt/global/core/workflow/type/edge'; import type { RuntimeEdgeItemType } from '@fastgpt/global/core/workflow/type/edge';
import type { RuntimeNodeItemType } from '@fastgpt/global/core/workflow/runtime/type'; import type { RuntimeNodeItemType } from '@fastgpt/global/core/workflow/runtime/type';
...@@ -75,7 +76,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -75,7 +76,7 @@ describe('getWorkflowResponseWrite', () => {
it('should not write when res is undefined', () => { it('should not write when res is undefined', () => {
const fn = getWorkflowResponseWrite({ detail: true, streamResponse: true }); const fn = getWorkflowResponseWrite({ detail: true, streamResponse: true });
fn({ event: SseResponseEventEnum.answer, data: { text: 'hi' } }); fn(workflowSseEvent.answerDelta('hi'));
// No error thrown // No error thrown
}); });
...@@ -84,7 +85,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -84,7 +85,7 @@ describe('getWorkflowResponseWrite', () => {
res.closed = true; res.closed = true;
vi.mocked(responseWrite).mockClear(); vi.mocked(responseWrite).mockClear();
const fn = getWorkflowResponseWrite({ res, detail: true, streamResponse: true }); const fn = getWorkflowResponseWrite({ res, detail: true, streamResponse: true });
fn({ event: SseResponseEventEnum.answer, data: { text: 'hi' } }); fn(workflowSseEvent.answerDelta('hi'));
expect(responseWrite).not.toHaveBeenCalled(); expect(responseWrite).not.toHaveBeenCalled();
}); });
...@@ -92,7 +93,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -92,7 +93,7 @@ describe('getWorkflowResponseWrite', () => {
const res = mockRes(); const res = mockRes();
vi.mocked(responseWrite).mockClear(); vi.mocked(responseWrite).mockClear();
const fn = getWorkflowResponseWrite({ res, detail: true, streamResponse: false }); const fn = getWorkflowResponseWrite({ res, detail: true, streamResponse: false });
fn({ event: SseResponseEventEnum.answer, data: { text: 'hi' } }); fn(workflowSseEvent.answerDelta('hi'));
expect(responseWrite).not.toHaveBeenCalled(); expect(responseWrite).not.toHaveBeenCalled();
}); });
...@@ -100,7 +101,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -100,7 +101,7 @@ describe('getWorkflowResponseWrite', () => {
const res = mockRes(); const res = mockRes();
vi.mocked(responseWrite).mockClear(); vi.mocked(responseWrite).mockClear();
const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true }); const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true });
fn({ event: SseResponseEventEnum.answer, data: { text: 'hi' } }); fn(workflowSseEvent.answerDelta('hi'));
expect(responseWrite).toHaveBeenCalled(); expect(responseWrite).toHaveBeenCalled();
}); });
...@@ -108,7 +109,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -108,7 +109,7 @@ describe('getWorkflowResponseWrite', () => {
const res = mockRes(); const res = mockRes();
vi.mocked(responseWrite).mockClear(); vi.mocked(responseWrite).mockClear();
const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true }); const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true });
fn({ event: SseResponseEventEnum.fastAnswer, data: { text: 'hi' } }); fn(workflowSseEvent.fastAnswerDelta('hi'));
expect(responseWrite).toHaveBeenCalled(); expect(responseWrite).toHaveBeenCalled();
}); });
...@@ -116,7 +117,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -116,7 +117,7 @@ describe('getWorkflowResponseWrite', () => {
const res = mockRes(); const res = mockRes();
vi.mocked(responseWrite).mockClear(); vi.mocked(responseWrite).mockClear();
const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true }); const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true });
fn({ event: SseResponseEventEnum.flowNodeStatus, data: {} }); fn(workflowSseEvent.flowNodeStatus('test node'));
expect(responseWrite).not.toHaveBeenCalled(); expect(responseWrite).not.toHaveBeenCalled();
}); });
...@@ -129,7 +130,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -129,7 +130,7 @@ describe('getWorkflowResponseWrite', () => {
streamResponse: true, streamResponse: true,
showNodeStatus: false showNodeStatus: false
}); });
fn({ event: SseResponseEventEnum.flowNodeStatus, data: {} }); fn(workflowSseEvent.flowNodeStatus('test node'));
expect(responseWrite).not.toHaveBeenCalled(); expect(responseWrite).not.toHaveBeenCalled();
}); });
...@@ -142,7 +143,15 @@ describe('getWorkflowResponseWrite', () => { ...@@ -142,7 +143,15 @@ describe('getWorkflowResponseWrite', () => {
streamResponse: true, streamResponse: true,
showNodeStatus: false showNodeStatus: false
}); });
fn({ event: SseResponseEventEnum.toolCall, data: {} }); fn(
workflowSseEvent.toolCall({
id: 'tool-call-id',
toolName: 'tool',
toolAvatar: '',
functionName: 'tool',
params: ''
})
);
expect(responseWrite).not.toHaveBeenCalled(); expect(responseWrite).not.toHaveBeenCalled();
}); });
...@@ -155,7 +164,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -155,7 +164,7 @@ describe('getWorkflowResponseWrite', () => {
streamResponse: true, streamResponse: true,
id: 'test-id' id: 'test-id'
}); });
fn({ id: 'rid', event: SseResponseEventEnum.answer, data: { text: 'hi' } }); fn(workflowSseEvent.answerDelta('hi', 'rid'));
expect(responseWrite).toHaveBeenCalledWith( expect(responseWrite).toHaveBeenCalledWith(
expect.objectContaining({ expect.objectContaining({
res, res,
...@@ -170,7 +179,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -170,7 +179,7 @@ describe('getWorkflowResponseWrite', () => {
const res = mockRes(); const res = mockRes();
vi.mocked(responseWrite).mockClear(); vi.mocked(responseWrite).mockClear();
const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true }); const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true });
fn({ event: SseResponseEventEnum.answer, data: { text: 'hi' } }); fn(workflowSseEvent.answerDelta('hi'));
expect(responseWrite).toHaveBeenCalledWith(expect.objectContaining({ event: undefined })); expect(responseWrite).toHaveBeenCalledWith(expect.objectContaining({ event: undefined }));
}); });
...@@ -178,7 +187,7 @@ describe('getWorkflowResponseWrite', () => { ...@@ -178,7 +187,7 @@ describe('getWorkflowResponseWrite', () => {
const res = mockRes(); const res = mockRes();
vi.mocked(responseWrite).mockClear(); vi.mocked(responseWrite).mockClear();
const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true }); const fn = getWorkflowResponseWrite({ res, detail: false, streamResponse: true });
fn({ event: SseResponseEventEnum.chatTitle, data: { title: 'Generated Title' } }); fn(workflowSseEvent.chatTitle('Generated Title'));
expect(responseWrite).toHaveBeenCalledWith( expect(responseWrite).toHaveBeenCalledWith(
expect.objectContaining({ expect.objectContaining({
event: SseResponseEventEnum.chatTitle event: SseResponseEventEnum.chatTitle
...@@ -224,15 +233,11 @@ describe('getWorkflowChildResponseWrite', () => { ...@@ -224,15 +233,11 @@ describe('getWorkflowChildResponseWrite', () => {
fn: mockFn as any fn: mockFn as any
}); });
expect(wrapped).toBeDefined(); expect(wrapped).toBeDefined();
wrapped!({ wrapped!(workflowSseEvent.answerDelta('hi'));
event: SseResponseEventEnum.answer,
data: { text: 'hi' }
});
expect(mockFn).toHaveBeenCalledWith( expect(mockFn).toHaveBeenCalledWith(
expect.objectContaining({ expect.objectContaining({
id: 'child-id', id: 'child-id',
event: SseResponseEventEnum.answer, event: SseResponseEventEnum.answer
data: { text: 'hi' }
}) })
); );
}); });
...@@ -244,16 +249,17 @@ describe('getWorkflowChildResponseWrite', () => { ...@@ -244,16 +249,17 @@ describe('getWorkflowChildResponseWrite', () => {
fn: mockFn as any fn: mockFn as any
}); });
wrapped!({ wrapped!(
id: 'tool-call-id', workflowSseEvent.toolParams({
event: SseResponseEventEnum.toolCall, id: 'tool-call-id',
data: { tool: { id: 'tool-call-id' } } params: 'delta'
}); })
);
expect(mockFn).toHaveBeenCalledWith( expect(mockFn).toHaveBeenCalledWith(
expect.objectContaining({ expect.objectContaining({
id: 'tool-call-id', id: 'tool-call-id',
event: SseResponseEventEnum.toolCall event: SseResponseEventEnum.toolParams
}) })
); );
}); });
......
...@@ -13,6 +13,7 @@ import { ...@@ -13,6 +13,7 @@ import {
initWorkflowSseResponse initWorkflowSseResponse
} from '@fastgpt/service/core/workflow/utils/streamResponseContext'; } from '@fastgpt/service/core/workflow/utils/streamResponseContext';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants';
import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { import {
ChatSourceTypeEnum, ChatSourceTypeEnum,
STREAM_RESUME_REQUEST_HEADER STREAM_RESUME_REQUEST_HEADER
...@@ -137,10 +138,7 @@ describe('createWorkflowStreamResponseContext', () => { ...@@ -137,10 +138,7 @@ describe('createWorkflowStreamResponseContext', () => {
responseId: chatId responseId: chatId
}); });
noMirrorContext.responseWrite({ noMirrorContext.responseWrite(workflowSseEvent.done(SseResponseEventEnum.answer));
event: SseResponseEventEnum.answer,
data: '[DONE]'
});
await noMirrorContext.flushResume(); await noMirrorContext.flushResume();
expect(redis.info).not.toHaveBeenCalled(); expect(redis.info).not.toHaveBeenCalled();
...@@ -158,10 +156,7 @@ describe('createWorkflowStreamResponseContext', () => { ...@@ -158,10 +156,7 @@ describe('createWorkflowStreamResponseContext', () => {
responseId: chatId responseId: chatId
}); });
mirrorContext.responseWrite({ mirrorContext.responseWrite(workflowSseEvent.done(SseResponseEventEnum.answer));
event: SseResponseEventEnum.answer,
data: '[DONE]'
});
await mirrorContext.flushResume(); await mirrorContext.flushResume();
expect(redis.info).toHaveBeenCalledTimes(1); expect(redis.info).toHaveBeenCalledTimes(1);
...@@ -194,10 +189,7 @@ describe('createWorkflowStreamResponseContext', () => { ...@@ -194,10 +189,7 @@ describe('createWorkflowStreamResponseContext', () => {
enableStreamResume: false enableStreamResume: false
}); });
context.responseWrite({ context.responseWrite(workflowSseEvent.done(SseResponseEventEnum.answer));
event: SseResponseEventEnum.answer,
data: '[DONE]'
});
expect(context).not.toHaveProperty('flushResume'); expect(context).not.toHaveProperty('flushResume');
expect(redis.info).not.toHaveBeenCalled(); expect(redis.info).not.toHaveBeenCalled();
......
...@@ -7,20 +7,19 @@ import type { ...@@ -7,20 +7,19 @@ import type {
SkillModuleResponseItemType SkillModuleResponseItemType
} from '@fastgpt/global/core/chat/type'; } from '@fastgpt/global/core/chat/type';
import type { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants'; import type { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants';
import type { WorkflowToolDeltaType } from '@fastgpt/global/core/workflow/runtime/sse';
import type { WorkflowInteractiveResponseType } from '@fastgpt/global/core/workflow/template/system/interactive/type'; import type { WorkflowInteractiveResponseType } from '@fastgpt/global/core/workflow/template/system/interactive/type';
import type { ChatAgentConfigFormDataType } from '@fastgpt/global/core/ai/auxiliaryGeneration/type'; import type { ChatAgentConfigFormDataType } from '@fastgpt/global/core/ai/auxiliaryGeneration/type';
import type { AuxiliaryGenerationEventEnum } from '@fastgpt/global/core/ai/auxiliaryGeneration/constants'; import type { AuxiliaryGenerationEventEnum } from '@fastgpt/global/core/ai/auxiliaryGeneration/constants';
import type { AgentPlanStatusType, AgentPlanType } from '@fastgpt/global/core/ai/agent/type'; import type { AgentPlanStatusType, AgentPlanType } from '@fastgpt/global/core/ai/agent/type';
export type generatingMessageProps = { type BaseGeneratingMessageProps = {
event: SseResponseEventEnum | AuxiliaryGenerationEventEnum;
responseValueId?: string; responseValueId?: string;
text?: string; text?: string;
reasoningText?: string; reasoningText?: string;
name?: string; name?: string;
status?: 'running' | 'finish'; status?: 'running' | 'finish';
tool?: ToolModuleResponseItemType;
interactive?: WorkflowInteractiveResponseType; interactive?: WorkflowInteractiveResponseType;
variables?: Record<string, any>; variables?: Record<string, any>;
nodeResponse?: ChatHistoryItemResType; nodeResponse?: ChatHistoryItemResType;
...@@ -38,6 +37,25 @@ export type generatingMessageProps = { ...@@ -38,6 +37,25 @@ export type generatingMessageProps = {
formData?: ChatAgentConfigFormDataType; formData?: ChatAgentConfigFormDataType;
}; };
type ToolStreamEvent =
| SseResponseEventEnum.toolCall
| SseResponseEventEnum.toolParams
| SseResponseEventEnum.toolResponse;
export type generatingMessageProps =
| (BaseGeneratingMessageProps & {
event: SseResponseEventEnum.toolCall;
tool?: ToolModuleResponseItemType;
})
| (BaseGeneratingMessageProps & {
event: SseResponseEventEnum.toolParams | SseResponseEventEnum.toolResponse;
tool?: WorkflowToolDeltaType;
})
| (BaseGeneratingMessageProps & {
event: Exclude<SseResponseEventEnum | AuxiliaryGenerationEventEnum, ToolStreamEvent>;
tool?: never;
});
export type StartChatFnProps = { export type StartChatFnProps = {
messages: ChatCompletionMessageParam[]; messages: ChatCompletionMessageParam[];
responseChatItemId?: string; responseChatItemId?: string;
......
...@@ -21,9 +21,9 @@ import { ...@@ -21,9 +21,9 @@ import {
getWorkflowEntryNodeIds, getWorkflowEntryNodeIds,
storeEdges2RuntimeEdges, storeEdges2RuntimeEdges,
rewriteNodeOutputByHistories, rewriteNodeOutputByHistories,
storeNodes2RuntimeNodes, storeNodes2RuntimeNodes
textAdaptGptResponse
} from '@fastgpt/global/core/workflow/runtime/utils'; } from '@fastgpt/global/core/workflow/runtime/utils';
import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { WORKFLOW_MAX_RUN_TIMES } from '@fastgpt/service/core/workflow/constants'; import { WORKFLOW_MAX_RUN_TIMES } from '@fastgpt/service/core/workflow/constants';
import { getWorkflowToolInputsFromStoreNodes } from '@fastgpt/global/core/app/tool/workflowTool/utils'; import { getWorkflowToolInputsFromStoreNodes } from '@fastgpt/global/core/app/tool/workflowTool/utils';
import { getChatItems } from '@fastgpt/service/core/chat/controller'; import { getChatItems } from '@fastgpt/service/core/chat/controller';
...@@ -262,17 +262,8 @@ async function handler(req: NextApiRequest, res: NextApiResponse) { ...@@ -262,17 +262,8 @@ async function handler(req: NextApiRequest, res: NextApiResponse) {
} }
}); });
streamResponseContext.responseWrite({ streamResponseContext.responseWrite(workflowSseEvent.answerStop());
event: SseResponseEventEnum.answer, streamResponseContext.responseWrite(workflowSseEvent.done(SseResponseEventEnum.answer));
data: textAdaptGptResponse({
text: null,
finish_reason: 'stop'
})
});
streamResponseContext.responseWrite({
event: SseResponseEventEnum.answer,
data: '[DONE]'
});
const aiResponse: AIChatItemType & { dataId?: string } = { const aiResponse: AIChatItemType & { dataId?: string } = {
dataId: responseChatItemId, dataId: responseChatItemId,
......
...@@ -13,9 +13,9 @@ import { ...@@ -13,9 +13,9 @@ import {
getMaxHistoryLimitFromNodes, getMaxHistoryLimitFromNodes,
storeEdges2RuntimeEdges, storeEdges2RuntimeEdges,
storeNodes2RuntimeNodes, storeNodes2RuntimeNodes,
textAdaptGptResponse,
getLastInteractiveValue getLastInteractiveValue
} from '@fastgpt/global/core/workflow/runtime/utils'; } from '@fastgpt/global/core/workflow/runtime/utils';
import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { GPTMessages2Chats, chatValue2RuntimePrompt } from '@fastgpt/global/core/chat/adapt'; import { GPTMessages2Chats, chatValue2RuntimePrompt } from '@fastgpt/global/core/chat/adapt';
import { getChatItems } from '@fastgpt/service/core/chat/controller'; import { getChatItems } from '@fastgpt/service/core/chat/controller';
import { import {
...@@ -38,7 +38,6 @@ import { ...@@ -38,7 +38,6 @@ import {
} from '@fastgpt/global/core/chat/utils'; } from '@fastgpt/global/core/chat/utils';
import { updateApiKeyUsage } from '@fastgpt/service/support/openapi/tools'; import { updateApiKeyUsage } from '@fastgpt/service/support/openapi/tools';
import { getRunningUserInfoByTmbId } from '@fastgpt/service/support/user/team/utils'; import { getRunningUserInfoByTmbId } from '@fastgpt/service/support/user/team/utils';
import { AuthUserTypeEnum } from '@fastgpt/global/support/permission/constant';
import { MongoApp } from '@fastgpt/service/core/app/schema'; import { MongoApp } from '@fastgpt/service/core/app/schema';
import { type AuthOutLinkChatProps } from '@fastgpt/global/support/outLink/api'; import { type AuthOutLinkChatProps } from '@fastgpt/global/support/outLink/api';
import { MongoChat } from '@fastgpt/service/core/chat/chatSchema'; import { MongoChat } from '@fastgpt/service/core/chat/chatSchema';
...@@ -70,6 +69,7 @@ import { ...@@ -70,6 +69,7 @@ import {
filterWorkflowFinalResponseData, filterWorkflowFinalResponseData,
getWorkflowFinalResponseData getWorkflowFinalResponseData
} from '@/service/core/workflow/nodeResponse'; } from '@/service/core/workflow/nodeResponse';
import { formatCompletionResponseContent } from '@/service/core/chat/utils';
import { import {
createWorkflowStreamResponseContext, createWorkflowStreamResponseContext,
type WorkflowStreamResponseContext type WorkflowStreamResponseContext
...@@ -439,56 +439,27 @@ async function handler(req: NextApiRequest, res: NextApiResponse) { ...@@ -439,56 +439,27 @@ async function handler(req: NextApiRequest, res: NextApiResponse) {
await titleSender.send(); await titleSender.send();
titleSender.close(); titleSender.close();
workflowResponseWrite({ workflowResponseWrite(workflowSseEvent.answerStop());
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
text: null,
finish_reason: 'stop'
})
});
// 特殊输配(data 不是{}) // 特殊输配(data 不是{})
if (detail) { if (detail) {
workflowResponseWrite({ workflowResponseWrite(
event: SseResponseEventEnum.flowResponses, workflowSseEvent.raw({
data: JSON.stringify(feResponseData) event: SseResponseEventEnum.flowResponses,
}); data: JSON.stringify(feResponseData)
})
);
} }
workflowResponseWrite({ workflowResponseWrite(
event: detail ? SseResponseEventEnum.answer : undefined, workflowSseEvent.done(detail ? SseResponseEventEnum.answer : undefined)
data: '[DONE]' );
});
} else { } else {
const generatedTitle = await titleSender.send(); const generatedTitle = await titleSender.send();
const formatResponseContent = removeAIResponseCite(assistantResponses, retainDatasetCite); const formatResponseContent = removeAIResponseCite(assistantResponses, retainDatasetCite);
const formattdResponse = (() => { const formattdResponse = formatCompletionResponseContent({
if (formatResponseContent.length === 0) responseContent: formatResponseContent,
return { detail
reasoning: '', });
content: ''
};
if (formatResponseContent.length === 1) {
return {
reasoning: formatResponseContent[0].reasoning?.content,
content: formatResponseContent[0].text?.content
};
}
if (!detail) {
return {
reasoning: formatResponseContent
.map((item) => item?.reasoning?.content)
.filter(Boolean)
.join('\n'),
content: formatResponseContent
.map((item) => item?.text?.content)
.filter(Boolean)
.join('\n')
};
}
return formatResponseContent;
})();
const error = const error =
nodeResponseSummary?.lastError || nodeResponseSummary?.lastError ||
...@@ -507,30 +478,7 @@ async function handler(req: NextApiRequest, res: NextApiResponse) { ...@@ -507,30 +478,7 @@ async function handler(req: NextApiRequest, res: NextApiResponse) {
message: { message: {
role: 'assistant', role: 'assistant',
...(Array.isArray(formattdResponse) ...(Array.isArray(formattdResponse)
? { ? { content: formattdResponse }
content: formattdResponse.map((item) => {
// Add value type(适配旧版)
enum ChatItemValueTypeEnum {
text = 'text',
file = 'file',
tool = 'tool',
interactive = 'interactive',
reasoning = 'reasoning'
}
const type = (() => {
if (item.text) return ChatItemValueTypeEnum.text;
if ('file' in item) return ChatItemValueTypeEnum.file;
if ('tool' in item || 'tools' in item) return ChatItemValueTypeEnum.tool;
if ('interactive' in item) return ChatItemValueTypeEnum.interactive;
if ('reasoning' in item) return ChatItemValueTypeEnum.reasoning;
return ChatItemValueTypeEnum.text;
})();
return {
...item,
type
};
})
}
: { : {
content: formattdResponse.content, content: formattdResponse.content,
...(formattdResponse.reasoning && { ...(formattdResponse.reasoning && {
......
...@@ -13,9 +13,9 @@ import { ...@@ -13,9 +13,9 @@ import {
getMaxHistoryLimitFromNodes, getMaxHistoryLimitFromNodes,
storeEdges2RuntimeEdges, storeEdges2RuntimeEdges,
storeNodes2RuntimeNodes, storeNodes2RuntimeNodes,
textAdaptGptResponse,
getLastInteractiveValue getLastInteractiveValue
} from '@fastgpt/global/core/workflow/runtime/utils'; } from '@fastgpt/global/core/workflow/runtime/utils';
import { workflowSseEvent } from '@fastgpt/global/core/workflow/runtime/sse';
import { GPTMessages2Chats, chatValue2RuntimePrompt } from '@fastgpt/global/core/chat/adapt'; import { GPTMessages2Chats, chatValue2RuntimePrompt } from '@fastgpt/global/core/chat/adapt';
import { getChatItems } from '@fastgpt/service/core/chat/controller'; import { getChatItems } from '@fastgpt/service/core/chat/controller';
import { import {
...@@ -37,7 +37,6 @@ import { ...@@ -37,7 +37,6 @@ import {
} from '@fastgpt/global/core/chat/utils'; } from '@fastgpt/global/core/chat/utils';
import { updateApiKeyUsage } from '@fastgpt/service/support/openapi/tools'; import { updateApiKeyUsage } from '@fastgpt/service/support/openapi/tools';
import { getRunningUserInfoByTmbId } from '@fastgpt/service/support/user/team/utils'; import { getRunningUserInfoByTmbId } from '@fastgpt/service/support/user/team/utils';
import { AuthUserTypeEnum } from '@fastgpt/global/support/permission/constant';
import { MongoApp } from '@fastgpt/service/core/app/schema'; import { MongoApp } from '@fastgpt/service/core/app/schema';
import { type AuthOutLinkChatProps } from '@fastgpt/global/support/outLink/api'; import { type AuthOutLinkChatProps } from '@fastgpt/global/support/outLink/api';
import { MongoChat } from '@fastgpt/service/core/chat/chatSchema'; import { MongoChat } from '@fastgpt/service/core/chat/chatSchema';
...@@ -69,6 +68,7 @@ import { ...@@ -69,6 +68,7 @@ import {
filterWorkflowFinalResponseData, filterWorkflowFinalResponseData,
getWorkflowFinalResponseData getWorkflowFinalResponseData
} from '@/service/core/workflow/nodeResponse'; } from '@/service/core/workflow/nodeResponse';
import { formatCompletionResponseContent } from '@/service/core/chat/utils';
import { import {
createWorkflowStreamResponseContext, createWorkflowStreamResponseContext,
type WorkflowStreamResponseContext type WorkflowStreamResponseContext
...@@ -449,51 +449,20 @@ async function handler(req: NextApiRequest, res: NextApiResponse) { ...@@ -449,51 +449,20 @@ async function handler(req: NextApiRequest, res: NextApiResponse) {
await titleSender.send(); await titleSender.send();
titleSender.close(); titleSender.close();
streamResponseContext.responseWrite({ streamResponseContext.responseWrite(workflowSseEvent.answerStop());
event: SseResponseEventEnum.answer,
data: textAdaptGptResponse({
text: null,
finish_reason: 'stop'
})
});
streamResponseContext.responseWrite({ streamResponseContext.responseWrite(
event: detail ? SseResponseEventEnum.answer : undefined, workflowSseEvent.done(detail ? SseResponseEventEnum.answer : undefined)
data: '[DONE]' );
});
await streamResponseContext.flushResume(); await streamResponseContext.flushResume();
} else { } else {
const generatedTitle = await titleSender.send(); const generatedTitle = await titleSender.send();
const formatResponseContent = removeAIResponseCite(assistantResponses, retainDatasetCite); const formatResponseContent = removeAIResponseCite(assistantResponses, retainDatasetCite);
const formattdResponse = (() => { const formattdResponse = formatCompletionResponseContent({
if (formatResponseContent.length === 0) responseContent: formatResponseContent,
return { detail
reasoning: '', });
content: ''
};
if (formatResponseContent.length === 1) {
return {
reasoning: formatResponseContent[0].reasoning?.content,
content: formatResponseContent[0].text?.content
};
}
if (!detail) {
return {
reasoning: formatResponseContent
.map((item) => item?.reasoning?.content)
.filter(Boolean)
.join('\n'),
content: formatResponseContent
.map((item) => item?.text?.content)
.filter(Boolean)
.join('\n')
};
}
return formatResponseContent;
})();
res.json({ res.json({
...(detail ? { responseData: feResponseData, newVariables } : {}), ...(detail ? { responseData: feResponseData, newVariables } : {}),
......
import type { DatasetCiteItemType } from '@fastgpt/global/core/dataset/type'; import type { DatasetCiteItemType } from '@fastgpt/global/core/dataset/type';
import type { AIChatItemValueItemType } from '@fastgpt/global/core/chat/type';
import type { WorkflowInteractiveResponseType } from '@fastgpt/global/core/workflow/template/system/interactive/type';
export type PublicCompletionInteractive = Pick<WorkflowInteractiveResponseType, 'type' | 'params'>;
export type CompletionDetailValueItem = Omit<AIChatItemValueItemType, 'interactive'> & {
interactive?: PublicCompletionInteractive;
};
export type CompletionResponseContent =
| {
reasoning?: string;
content?: string;
}
| CompletionDetailValueItem[];
/**
* 格式化 completions 非流式响应的 choices[].message.content。
*
* 普通单条文本保持 OpenAI 兼容字符串;detail=true 且包含交互、工具、文件或多段 value 时,
* 返回按字段名区分的 value 数组,确保交互节点可从 choices 中被客户端识别。
*/
export const formatCompletionResponseContent = ({
responseContent,
detail
}: {
responseContent: AIChatItemValueItemType[];
detail: boolean;
}): CompletionResponseContent => {
const MAX_COMPLETION_INTERACTIVE_EXTRACT_DEPTH = 20;
const getChildInteractiveResponse = (
interactive: WorkflowInteractiveResponseType
): WorkflowInteractiveResponseType | undefined => {
const childrenResponse = (interactive.params as { childrenResponse?: unknown })
.childrenResponse;
if (!childrenResponse || typeof childrenResponse !== 'object') return;
if (!('type' in childrenResponse) || !('params' in childrenResponse)) return;
return childrenResponse as WorkflowInteractiveResponseType;
};
const omitChildrenResponse = (
params: WorkflowInteractiveResponseType['params']
): WorkflowInteractiveResponseType['params'] => {
const publicParams = { ...(params as Record<string, unknown>) };
delete publicParams.childrenResponse;
return publicParams as WorkflowInteractiveResponseType['params'];
};
/**
* workflow 暂停恢复依赖 entryNodeIds、memoryEdges、nodeOutputs 等运行态字段;这些字段会保存到
* chat history,外部接口只需要返回当前交互的展示配置。children/loop/tool 包装节点会继续向下取
* 真正面向用户的交互节点,并限制最大深度,避免异常循环结构导致死循环。
*/
const formatPublicCompletionInteractive = (
interactive: WorkflowInteractiveResponseType
): PublicCompletionInteractive => {
let current = interactive;
let depth = 0;
const visited = new WeakSet<object>();
while (depth < MAX_COMPLETION_INTERACTIVE_EXTRACT_DEPTH) {
if (visited.has(current)) break;
visited.add(current);
const childrenResponse = getChildInteractiveResponse(current);
if (!childrenResponse) break;
current = childrenResponse;
depth++;
}
return {
type: current.type,
params: omitChildrenResponse(current.params)
};
};
const formatDetailValueList = (responseContent: AIChatItemValueItemType[]) =>
responseContent.map((item) => ({
...item,
...(item.interactive && {
interactive: formatPublicCompletionInteractive(item.interactive)
})
}));
const shouldReturnDetailValueList = ({
responseContent,
detail
}: {
responseContent: AIChatItemValueItemType[];
detail: boolean;
}) => {
if (!detail) return false;
if (responseContent.length > 1) return true;
const [item] = responseContent;
if (!item) return false;
return !!(item.interactive || item.tool || item.tools || 'file' in item);
};
if (responseContent.length === 0) {
return {
reasoning: '',
content: ''
};
}
if (shouldReturnDetailValueList({ responseContent, detail })) {
return formatDetailValueList(responseContent);
}
if (responseContent.length === 1) {
return {
reasoning: responseContent[0].reasoning?.content,
content: responseContent[0].text?.content
};
}
return {
reasoning: responseContent
.map((item) => item?.reasoning?.content)
.filter(Boolean)
.join('\n'),
content: responseContent
.map((item) => item?.text?.content)
.filter(Boolean)
.join('\n')
};
};
// 获取对话时间时,引用的内容 // 获取对话时间时,引用的内容
export function processChatTimeFilter( export function processChatTimeFilter(
......
...@@ -23,13 +23,9 @@ import { formatTime2YMDHMW } from '@fastgpt/global/common/string/time'; ...@@ -23,13 +23,9 @@ import { formatTime2YMDHMW } from '@fastgpt/global/common/string/time';
import { getWebReqUrl } from '@fastgpt/web/common/system/utils'; import { getWebReqUrl } from '@fastgpt/web/common/system/utils';
import type { OnOptimizePromptProps } from '@/components/common/PromptEditor/OptimizerPopover'; import type { OnOptimizePromptProps } from '@/components/common/PromptEditor/OptimizerPopover';
import type { OnOptimizeCodeProps } from '@/pageComponents/app/detail/WorkflowComponents/Flow/nodes/NodeCode/Copilot'; import type { OnOptimizeCodeProps } from '@/pageComponents/app/detail/WorkflowComponents/Flow/nodes/NodeCode/Copilot';
import type { import type { WorkflowSsePayloadMap } from '@fastgpt/global/core/workflow/runtime/sse';
ToolModuleResponseItemType,
SkillModuleResponseItemType
} from '@fastgpt/global/core/chat/type';
import type { ChatAgentConfigFormDataType } from '@fastgpt/global/core/ai/auxiliaryGeneration/type'; import type { ChatAgentConfigFormDataType } from '@fastgpt/global/core/ai/auxiliaryGeneration/type';
import { AuxiliaryGenerationEventEnum } from '@fastgpt/global/core/ai/auxiliaryGeneration/constants'; import { AuxiliaryGenerationEventEnum } from '@fastgpt/global/core/ai/auxiliaryGeneration/constants';
import type { AgentPlanStatusType, AgentPlanType } from '@fastgpt/global/core/ai/agent/type';
import type { StreamNoNeedToBeResumeType } from '@fastgpt/global/openapi/core/ai/api'; import type { StreamNoNeedToBeResumeType } from '@fastgpt/global/openapi/core/ai/api';
type StreamFetchProps = { type StreamFetchProps = {
...@@ -69,45 +65,31 @@ const shouldSendStreamResumeHeader = (url: string) => ...@@ -69,45 +65,31 @@ const shouldSendStreamResumeHeader = (url: string) =>
type CommonResponseType = { type CommonResponseType = {
responseValueId?: string; responseValueId?: string;
}; };
type ResponseQueueItemType = CommonResponseType & type WorkflowQueueEvent =
( | SseResponseEventEnum.interactive
| { | SseResponseEventEnum.toolCall
event: SseResponseEventEnum.fastAnswer | SseResponseEventEnum.answer; | SseResponseEventEnum.toolParams
text?: string; | SseResponseEventEnum.toolResponse
reasoningText?: string; | SseResponseEventEnum.plan
} | SseResponseEventEnum.planStatus
| { | SseResponseEventEnum.skillCall
event: SseResponseEventEnum.interactive; | SseResponseEventEnum.chatTitle;
[key: string]: any; type WorkflowQueueItem<Event extends WorkflowQueueEvent> = Event extends WorkflowQueueEvent
} ? CommonResponseType & { event: Event } & WorkflowSsePayloadMap[Event]
| { : never;
event: type AnswerQueueItem = CommonResponseType & {
| SseResponseEventEnum.toolCall event: SseResponseEventEnum.fastAnswer | SseResponseEventEnum.answer;
| SseResponseEventEnum.toolParams text?: string;
| SseResponseEventEnum.toolResponse; reasoningText?: string;
tool: ToolModuleResponseItemType; };
} type ChatAgentConfigQueueItem = CommonResponseType & {
| { event: AuxiliaryGenerationEventEnum.chatAgentConfig;
event: AuxiliaryGenerationEventEnum.chatAgentConfig; data: ChatAgentConfigFormDataType;
data: ChatAgentConfigFormDataType; };
} type ResponseQueueItemType =
| { | AnswerQueueItem
event: SseResponseEventEnum.plan; | WorkflowQueueItem<WorkflowQueueEvent>
plan: AgentPlanType; | ChatAgentConfigQueueItem;
}
| {
event: SseResponseEventEnum.planStatus;
planStatus: AgentPlanStatusType;
}
| {
event: SseResponseEventEnum.skillCall;
skill: SkillModuleResponseItemType;
}
| {
event: SseResponseEventEnum.chatTitle;
title: string;
}
);
const STREAM_TYPING_QUEUE_COUNT_WHILE_STREAMING = 1; const STREAM_TYPING_QUEUE_COUNT_WHILE_STREAMING = 1;
const STREAM_TYPING_QUEUE_COUNT_AFTER_FINISH = 20; const STREAM_TYPING_QUEUE_COUNT_AFTER_FINISH = 20;
......
import { describe, expect, it } from 'vitest';
import { formatCompletionResponseContent } from '@/service/core/chat/utils';
import type { AIChatItemValueItemType } from '@fastgpt/global/core/chat/type';
describe('formatCompletionResponseContent', () => {
it('keeps plain single text response OpenAI compatible', () => {
const result = formatCompletionResponseContent({
detail: true,
responseContent: [
{
text: {
content: 'hello'
}
}
]
});
expect(result).toEqual({
reasoning: undefined,
content: 'hello'
});
});
it('returns single interactive response as field-keyed content array when detail is enabled', () => {
const interactive = {
type: 'userSelect',
params: {
description: '请选择',
userSelectOptions: [{ label: 'A', value: 'A' }]
}
} as AIChatItemValueItemType['interactive'];
const result = formatCompletionResponseContent({
detail: true,
responseContent: [
{
interactive
}
]
});
expect(result).toEqual([
{
interactive
}
]);
});
it('returns only public interactive fields and extracts deepest child interactive', () => {
const interactive = {
type: 'childrenInteractive',
entryNodeIds: ['parent'],
interactiveId: 'internal-interactive-id',
nodeResponseId: 'internal-node-response-id',
memoryEdges: [{ id: 'edge' }],
nodeOutputs: [{ id: 'output' }],
params: {
childrenResponse: {
type: 'userInput',
entryNodeIds: ['child'],
interactiveId: 'child-interactive-id',
memoryEdges: [{ id: 'child-edge' }],
nodeOutputs: [{ id: 'child-output' }],
params: {
description: '填写信息',
inputForm: []
}
}
}
} as unknown as AIChatItemValueItemType['interactive'];
const result = formatCompletionResponseContent({
detail: true,
responseContent: [
{
interactive
}
]
});
expect(result).toEqual([
{
interactive: {
type: 'userInput',
params: {
description: '填写信息',
inputForm: []
}
}
}
]);
});
it('stops extracting cyclic child interactive and strips childrenResponse', () => {
const interactive = {
type: 'childrenInteractive',
params: {}
} as Record<string, any>;
interactive.params.childrenResponse = interactive;
const result = formatCompletionResponseContent({
detail: true,
responseContent: [
{
interactive: interactive as AIChatItemValueItemType['interactive']
}
]
});
expect(result).toEqual([
{
interactive: {
type: 'childrenInteractive',
params: {}
}
}
]);
});
it('joins multiple text and reasoning values when detail is disabled', () => {
const result = formatCompletionResponseContent({
detail: false,
responseContent: [
{
reasoning: {
content: 'think 1'
},
text: {
content: 'hello'
}
},
{
reasoning: {
content: 'think 2'
},
text: {
content: 'world'
}
}
]
});
expect(result).toEqual({
reasoning: 'think 1\nthink 2',
content: 'hello\nworld'
});
});
});
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or sign in to comment