Commit 9f383002 by Archer Committed by GitHub

perf: node store (#7112)

* perf: node store

* revert

* fix: loop

* doc

* feat: source

* perf: store chat

* perf: optimize chat node response preview reads

* perf: speed up persisted node response tree compose

* doc

* perf: depth

* perf: chat anination
parent 4fa14557
# Workflow NodeResponse 流式持久化与平铺存储
日期:2026-05-30
基线:`/Volumes/code/FastGPT-worktrees/upstream-main`
## 最终摘要
本轮工作将 workflow 运行期的 `nodeResponses/responseData` 从“内存累计,结束后一次保存”调整为“节点完成后通过请求级 writer 分批写入数据库”。workflow 子流程、loop、parallel、toolcall 等大子树默认在 `chat_item_responses` 中平铺保存,通过 `data.id/data.parentId` 还原父子关系;少数可控、小规模且语义上绑定当前节点的 child 详情允许继续内联保留。
客户端详情结构保持稳定:接口返回和流式展示仍是嵌套结构,正式 child 字段继续使用既有 `childrenResponses`;旧 detail 字段 `pluginDetail/toolDetail/loopDetail/loopRunDetail/parallelDetail` 仅作为历史读取来源,不再作为新链路的通用写入结构。
## 已确认设计
- 复用 `chat_item_responses` 表,不新增版本字段。
- 每条 row 的 `data` 就是完整的当前节点 `nodeResponse` 数据;父子识别只使用 `data.id``data.parentId`
- 不新增 `sequence/depth/rootResponseId/runId/childRuntime`;展示顺序依赖同一个 writer 串行写入后的数据库插入顺序。
- 父节点记录 `childTotalPoints``childResponseCount`,统计完整 child 树;节点自身只保留当前节点运行时间。
- SSE `flowNodeResponse` 透出 `id/parentId`,客户端可在流式过程中合并父子节点。
- 新数据不再写入 `chat_items.responseData``saveChat` 只保存 AI 消息主体,引用、错误数和积分日志来自 writer 的 `nodeResponseSummary`
- 旧数据兼容保留:如果 `chat_items.responseData` 已存在,说明该条 AI 消息仍是历史内联详情,读取时优先使用它,不用 `chat_item_responses` 覆盖。
## 数据结构
row 核心字段:
```ts
type ChatItemResponseSchema = {
teamId: ObjectId;
appId: ObjectId;
chatId: string;
chatItemDataId: string;
// 完整的当前节点 nodeResponse 数据。
// data.id 是节点响应实例 ID;data.parentId 指向父节点 data.id。
// data.nodeId/moduleType/childTotalPoints/childResponseCount 都只保存在 data 内。
// 大子流程通过 parentId 平铺落库,读取时再拼回 childrenResponses。
// 少数可控、小规模的节点内联 child 可保留在 data 内。
data: ChatHistoryItemResType;
time: Date;
};
```
新增索引:
- `{ appId, chatId, chatItemDataId, _id }`,用于详情读取并按写入顺序排序
- `{ appId, chatId, chatItemDataId, data.id }`,unique
- `{ teamId, time: -1 }`
## 写入流程
实现位置:
- [nodeResponseStorage.ts](/Volumes/code/FastGPT-worktrees/upstream-main/packages/service/core/chat/nodeResponseStorage.ts)
- [dispatch/index.ts](/Volumes/code/FastGPT-worktrees/upstream-main/packages/service/core/workflow/dispatch/index.ts)
关键行为:
1. root chat v2 workflow 创建一个 `WorkflowNodeResponseWriter`,子 workflow、loop、parallel、plugin、runApp、toolcall 复用同一个 writer。
2. 节点执行前生成节点响应实例 ID 并写入 `data.id`;子 workflow 的 `nodeResponseParentId` 指向父节点或虚拟 wrapper 的 `data.id`,落库和 SSE 时表现为 `parentId`
3. 节点完成后调用 `record()` 写入轻量化后的 rows;大子流程 child 通过独立 row 平铺保存,父节点只保留统计信息。
4. 模块返回的内部 `nodeResponses` 如果没有显式 `parentId`,会默认挂到当前节点响应下;父节点缺少 child 统计时按这些 child 自动补 `childTotalPoints/childResponseCount`
5. writer 默认 `batchSize = 5`,buffer 满或 root close 时 flush。这里按 row 计数,大子流程会展开成多条 row;可控小 child 可能作为当前节点的内联详情随同一 row 写入。
6. flush 使用同一 Mongo transaction,先删除本批同 `data.id` 旧 row,再一次性 `create(rows[])``create` 设置 `ordered: true`,同一 flush 内全部成功或全部回滚。写入前不再用 `JSON.stringify` 预估体积,BSON 大小和不可序列化字段统一交给 Mongo 写入校验。
7. 普通写入失败重试 3 次;仍失败则将 rows 瘦身为节点名、类型、头像、运行时间、消耗统计、父子统计等关键字段后再写 1 次。
8. 瘦身写仍失败时丢弃本批 rows,记录日志并继续 workflow;失败 rows 不回退到内存,也不交给 `saveChat` 兜底。
9. 详情 rows 被丢弃时,writer 已累计的 summary 仍保留,`saveChat` 仍可保存引用、错误数和积分日志。
10. `record()` 返回规范化后的响应。大子流程 child 已经通过独立 row 交给 writer,调用链可释放原始完整 `nodeResponse`;可控小 child 仍可作为当前节点内联详情保留。
11. parallel retry 通过 `deleteResponses()` 清理失败 attempt 及其 descendants。
12. 新回合使用 replace;interactive resume 使用 append。replace 的旧详情删除延迟到第一次 flush transaction 中执行,避免新写失败先清空旧详情。
## 读取规则
实现位置:
- [controller.ts](/Volumes/code/FastGPT-worktrees/upstream-main/packages/service/core/chat/controller.ts)
- [getResData.ts](/Volumes/code/FastGPT-worktrees/upstream-main/projects/app/src/pages/api/core/chat/record/getResData.ts)
- [utils.ts](/Volumes/code/FastGPT-worktrees/upstream-main/packages/global/core/chat/utils.ts)
读取规则:
- 对旧 AI 消息,如果 `chat_items.responseData` 已存在,直接返回内联详情,避免被独立表空结果覆盖。
- 对新 AI 消息,`chat_items.responseData` 不存在,详情接口和 chat history 从 `chat_item_responses` 读取 node response rows。
- rows 按 `_id: 1` 读取,通过 `data.parentId` 拼回 `childrenResponses`
- child row 早于 parent row 写入也能还原。
- ResponseTags、WholeResponseModal、SSE resume 都递归读取 `childrenResponses`,同时读取旧 detail 字段。
## 已完成改动
- 新增 `nodeResponseStorage.ts`,封装 flat row 生成、详情拼接、writer、summary、失败降级与 retry 清理。
- `MongoChatItemResponse` schema 收敛为所有节点详情都在 `data` 内,唯一性由同一 chat item 下的 `data.id` 保证。
- workflow dispatch 接入请求级 writer,root runtime 不再累计完整 `nodeResponses`
- loop/parallel/runApp/plugin/toolcall/agent 子流程复用 writer,并补齐父节点 child 统计。
- 修正子流程返回语义:持久化和 SSE 使用完整轻量响应,workflow 队列内部避免同一节点同时从 `responseData/nodeResponses` 重复进入 `flowResponses`
- `saveChat` 改为只使用 `nodeResponseSummary`,并显式丢弃误传的内联 `responseData`
- `pushChatLog` 只从持久化 rows 读取详情计算 `responseTime`
- `MongoChatItem` 的 deprecated `responseData` schema path 改为 `default: undefined`,避免新 chat item 自动带空数组。
- 前端流式合并支持 child 先到、parent 后到、重复 `id` 更新,以及详情弹窗递归展示。
## 验证覆盖
重点覆盖:
- rows 写入、child stats、受控 child/detail 字段保留、quote q/a 裁剪。
- 详情读取、旧 `responseData` 字段存在时优先保留、child 先到、append 合并。
- writer 串行顺序、批量 ordered create、session 复用、重复 `data.id` 覆盖、3 次重试、瘦身写入、最终失败丢弃。
- 写入失败丢弃 rows 后 summary 仍保留,`saveChat` 不依赖失败 rows。
- root workflow writer close、loop/parallel/toolcall 子响应写入、parallel retry 清理。
- 真实 `runWorkflow` 集成测试覆盖 loop 子 workflow 平铺落库、模块内部 child 自动挂当前节点、中途 batch flush 后的 `childrenResponses` 拼接、SSE `id/parentId`
- 前端 SSE parent/child 合并、child 先到、重复更新、ResponseTags/WholeResponseModal 递归读取。
最近验证命令:
- `pnpm --filter @fastgpt/service test test/core/chat/nodeResponseStorage.test.ts test/core/workflow/dispatch/loopRun/runLoopRun.test.ts`:2 个文件、39 个测试通过。
- `pnpm --filter @fastgpt/service test test/core/chat/controller.test.ts`:1 个文件、56 个测试通过。
- `pnpm --filter @fastgpt/app test test/api/core/chat/record/getResData.test.ts`:1 个文件、4 个测试通过。
- `pnpm --filter @fastgpt/app test test/components/core/chat/ChatContainer/ChatBox/resume.test.ts test/components/core/chat/ChatContainer/ChatBox/utils.test.ts test/global/core/chat/utils.test.ts`:3 个文件、36 个测试通过。
- `pnpm --filter @fastgpt/global test test/core/chat/chatUtils.test.ts`:1 个文件、53 个测试通过。
- `pnpm --filter @fastgpt/service test`:170 个文件通过、2 个文件跳过;2554 个测试通过、29 个跳过。
- `pnpm --filter @fastgpt/global test`:72 个文件、1626 个测试通过。
- `pnpm --filter @fastgpt/app test`:103 个文件、796 个测试通过。
- `pnpm --filter @fastgpt/app run typecheck`:通过。
- `pnpm exec eslint <changed ts/tsx files>`:通过。
- `git diff --check`
## 生产检查结论
- 当前实现已满足本轮确认的核心约束:单 writer 保序、大子流程平铺落库、允许部分可控节点保留小规模内联 child、按批次批量写入、失败重试与瘦身降级、失败后释放详情 rows、`saveChat` 仅依赖 summary、新数据不再落 `chat_items.responseData`、接口和前端继续返回 `childrenResponses` 嵌套结构。
- 仍建议上线后对真实大工作流采集 heap/rss、Mongo 写入耗时、接口耗时和 fallback 日志数量,用于确认生产数据分布下的默认 `batchSize` 是否需要调优。
......@@ -60,109 +60,3 @@ description: 'FastGPT V4.15.0-beta4 更新说明'
1. 插件服务从旧 `runtime` 结构调整为 pnpm workspace monorepo,拆分为 HTTP 服务入口、领域模型、用例、API adapter、基础设施、SDK 和 CLI。
2. 将 app API 接口全部用 zod schema 编写并生成文档。
3. 及时处理 worker 内图片,不再存留 base64,降低内存消耗。
## FAQ
如果出现 Mongo sync index error,并且是提示:
```shell
modelName: 'chat_item_responses',
error: MongoServerError: Index build failed: a45ca5d9-8429-4602-80bc-312a22b92687: Collection fastgpt.chat_item_responses ( f2a417fe-6e7c-4ec2-8947-d8003755462f ) :: caused by :: E11000 duplicate key error collection: fastgpt.chat_item_responses index: appId_1_chatId_1_chatItemDataId_1_data.id_1 dup key: { appId: ObjectId('69c9110790de745f57e93919'), chatId: "jzt1JCcaoy4kC67hdhQlzUob", chatItemDataId: "xa0f1AMFr9CWj9uVC9TqGtc0", data.id: "jTUTS9j67N4CyZ17" }
```
可以进入 mongo shell,执行以下命令来清楚重复的数据。这里仅删除对话中每个节点的运行详情结果,无风险,可能是过去某次出现异常导致重复。
```shell
use fastgpt
const dryRun = true; // 确认后改成 false
const batchSize = 500;
const col = db.chat_item_responses;
let duplicateGroups = 0;
let duplicateDocs = 0;
let removableDocs = 0;
let deletedDocs = 0;
let pendingIds = [];
function flushDeletes() {
if (pendingIds.length === 0) return;
const ids = pendingIds;
pendingIds = [];
if (dryRun) {
print(`[dry-run] would delete ${ids.length} old docs`);
return;
}
const res = col.deleteMany({ _id: { $in: ids } });
deletedDocs += res.deletedCount;
print(`[delete] deleted ${res.deletedCount} old docs, total=${deletedDocs}`);
}
const cursor = col.aggregate(
[
{
$group: {
_id: {
appId: "$appId",
chatId: "$chatId",
chatItemDataId: "$chatItemDataId",
dataId: "$data.id"
},
count: { $sum: 1 }
}
},
{ $match: { count: { $gt: 1 } } },
{ $sort: { count: -1 } }
],
{ allowDiskUse: true }
);
while (cursor.hasNext()) {
const group = cursor.next();
duplicateGroups += 1;
duplicateDocs += group.count;
removableDocs += group.count - 1;
const key = group._id;
const filter = {
appId: key.appId,
chatId: key.chatId,
chatItemDataId: key.chatItemDataId,
"data.id": key.dataId
};
const docs = col
.find(filter, { _id: 1, time: 1 })
.sort({ _id: -1 })
.toArray();
const keepId = docs[0]?._id;
const removeIds = docs.slice(1).map((doc) => doc._id);
printjson({
group: duplicateGroups,
count: group.count,
keepId,
removeIds
});
pendingIds.push(...removeIds);
if (pendingIds.length >= batchSize) {
flushDeletes();
}
}
flushDeletes();
printjson({
dryRun,
duplicateGroups,
duplicateDocs,
removableDocs,
deletedDocs
});
```
......@@ -6,6 +6,9 @@ description: 'FastGPT V4.15.0-beta5 更新说明'
## 🚀 新增内容
1. HTTP 节点支持配置忽略 TLS 证书校验,适用于调用使用自签名证书或内部证书的 HTTPS 服务。
2. 支持目录深度环境变量,避免无限嵌套目录。
3. 对话框支持快速滚动到底部按键。
4. 参考 Lobe UI 优化流输出动效。
## ⚙️ 优化
......
......@@ -149,8 +149,8 @@
"content/plugin/model-presets.mdx": "2026-06-04T16:10:15+08:00",
"content/plugin/system-tool-development.en.mdx": "2026-06-09T16:03:58+08:00",
"content/plugin/system-tool-development.mdx": "2026-06-09T16:03:58+08:00",
"content/self-host/config/env.en.mdx": "2026-06-13T15:48:04+08:00",
"content/self-host/config/env.mdx": "2026-06-13T15:48:04+08:00",
"content/self-host/config/env.en.mdx": "2026-06-14T00:12:11+08:00",
"content/self-host/config/env.mdx": "2026-06-14T00:12:11+08:00",
"content/self-host/config/json.en.mdx": "2026-05-25T11:21:30+08:00",
"content/self-host/config/json.mdx": "2026-05-25T11:21:30+08:00",
"content/self-host/config/model/intro.en.mdx": "2026-06-04T16:10:15+08:00",
......@@ -275,9 +275,9 @@
"content/self-host/upgrading/4-15/41503.en.mdx": "2026-05-28T16:21:09+08:00",
"content/self-host/upgrading/4-15/41503.mdx": "2026-05-28T16:21:09+08:00",
"content/self-host/upgrading/4-15/41504.en.mdx": "2026-06-10T19:02:59+08:00",
"content/self-host/upgrading/4-15/41504.mdx": "2026-06-11T17:42:15+08:00",
"content/self-host/upgrading/4-15/41504.mdx": "2026-06-14T22:25:36+08:00",
"content/self-host/upgrading/4-15/41505.en.mdx": "2026-06-12T20:47:04+08:00",
"content/self-host/upgrading/4-15/41505.mdx": "2026-06-13T22:52:27+08:00",
"content/self-host/upgrading/4-15/41505.mdx": "2026-06-15T21:58:39+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/41.en.mdx": "2026-04-26T21:08:47+08:00",
......
......@@ -191,6 +191,62 @@ export const mergeAssistantFieldMessages = (messages: ChatCompletionMessageParam
return mergedMessages;
};
const isPureTextAiValue = (item: AIChatItemValueItemType) =>
!!item.text &&
!item.id &&
!item.planId &&
!item.reasoning &&
!item.tools &&
!item.skills &&
!item.interactive &&
!item.plan &&
!item.planStatus &&
!item.agentPlanUpdate &&
!item.agentAsk &&
!item.agentStopGate &&
!item.contextCheckpoint &&
!item.tool &&
!item.hideReason &&
!item.hideInUI;
/**
* 规整 AI chat value。
*
* 运行期可能把普通回答拆成很多连续 text value,甚至连续追加空 text 占位。这里在适配层
* 归一化这些纯文本片段:非空纯文本连续时合并,纯空 text 占位直接丢弃;如果最终没有
* 任何可保存 value,再补一个空 text。带 reasoning、interactive、tool、plan 等语义字段
* 的 value 保持独立边界。
*/
export const normalizeAIChatValue = (values: AIChatItemValueItemType[]) => {
const result: AIChatItemValueItemType[] = [];
values.forEach((item) => {
if (!isPureTextAiValue(item)) {
result.push(item);
return;
}
const text = item.text?.content || '';
if (!text) return;
const lastItem = result[result.length - 1];
if (lastItem && isPureTextAiValue(lastItem)) {
lastItem.text!.content += text;
return;
}
result.push(item);
});
if (result.length === 0) {
result.push({
text: { content: '' }
});
}
return result;
};
/**
* 将 FastGPT 内部 ChatItem 历史转换为 GPT request messages。
*
......
......@@ -187,7 +187,7 @@ export const DispatchNodeResponseSchema = z
customOutputs: z.record(z.string(), z.any()).optional().meta({ description: '自定义输出' }),
nodeInputs: z.record(z.string(), z.any()).optional().meta({ description: '节点输入' }),
nodeOutputs: z.record(z.string(), z.any()).optional().meta({ description: '节点输出' }),
mergeSignId: z.string().optional().meta({ description: '合并签名 ID' }),
mergeSignId: z.string().optional().meta({ description: '旧版合并签名 ID', deprecated: true }),
parentId: z.string().optional().meta({ description: '父节点响应实例 ID' }),
// bill
......
......@@ -8,6 +8,7 @@ import { ChatCompletionMessageParamSchema } from '../../../../ai/llm/type';
export const InteractiveBasicTypeSchema = z.object({
entryNodeIds: z.array(z.string()),
nodeResponseId: z.string().optional(),
memoryEdges: z.array(RuntimeEdgeItemTypeSchema),
nodeOutputs: z.array(NodeOutputItemSchema),
skipNodeQueue: z
......@@ -19,6 +20,7 @@ export type InteractiveBasicType = z.infer<typeof InteractiveBasicTypeSchema>;
const InteractiveNodeTypeSchema = z.object({
entryNodeIds: z.array(z.string()).optional(),
nodeResponseId: z.string().optional(),
memoryEdges: z.array(RuntimeEdgeItemTypeSchema).optional(),
nodeOutputs: z.array(NodeOutputItemSchema).optional()
});
......@@ -47,7 +49,8 @@ export const ToolCallChildrenInteractiveSchema = z.object({
})
})
});
export type ToolCallChildrenInteractive = z.infer<typeof ToolCallChildrenInteractiveSchema>;
export type ToolCallChildrenInteractive = InteractiveNodeType &
z.infer<typeof ToolCallChildrenInteractiveSchema>;
// Loop bode
export const LoopInteractiveSchema = z.object({
......
......@@ -7,6 +7,16 @@ import { ChatItemMiniSchema } from '../../../../core/chat/type';
import { AppTTSConfigTypeSchema } from '../../../../core/app/type';
import { GetChatTypeEnum } from '../../../../core/chat/constants';
const QueryStringArraySchema = z
.union([z.string(), z.array(z.string())])
.optional()
.transform((val) => {
if (!val) return undefined;
const values = Array.isArray(val) ? val : val.split(',');
return values.map((item) => item.trim()).filter(Boolean);
});
/* ============================================================================
* API: 获取对话响应详细数据
* Route: GET /api/core/chat/record/getResData
......@@ -35,6 +45,10 @@ export const DeleteChatRecordBodySchema = OutLinkChatAuthSchema.extend({
example: 'content123',
description: '要删除的消息 ID'
}),
contentIds: QueryStringArraySchema.meta({
example: ['content123', 'content456'],
description: '要删除的消息 ID 列表'
}),
delFile: z.coerce.boolean().optional().meta({
example: false,
description: '是否同时删除关联文件'
......
......@@ -115,13 +115,8 @@ export const FastLoginBodySchema = z.object({
export type FastLoginBodyType = z.infer<typeof FastLoginBodySchema>;
// ===== WeChat Login Result =====
export const WxLoginBodySchema = z.object({
inviterId: z.string().optional().meta({ description: '邀请人 ID' }),
export const WxLoginBodySchema = TrackRegisterParamsSchema.extend({
code: z.string().meta({ description: '微信登录 Code' }),
bd_vid: z.string().optional(),
msclkid: z.string().optional(),
fastgpt_sem: z.string().optional(),
sourceDomain: z.string().optional(),
language: LanguageSchema.optional().meta({ description: '语言' })
});
export type WxLoginBodyType = z.infer<typeof WxLoginBodySchema>;
......
......@@ -7,14 +7,18 @@ export const ShortUrlSchema = z.object({
});
export type ShortUrlParams = z.infer<typeof ShortUrlSchema>;
export const FastGPT_SEM_Schema = ShortUrlSchema.extend({
keyword: z.string().optional(),
search: z.string().optional(),
source: z.string().optional(),
sourceDomain: z.string().optional()
});
export type FastGPTSemType = z.infer<typeof FastGPT_SEM_Schema>;
export const TrackRegisterParamsSchema = z.object({
inviterId: z.string().optional(),
bd_vid: z.string().optional(),
msclkid: z.string().optional(),
fastgpt_sem: ShortUrlSchema.extend({
keyword: z.string().optional(),
search: z.string().optional()
}).optional(),
sourceDomain: z.string().optional()
fastgpt_sem: FastGPT_SEM_Schema.optional()
});
export type TrackRegisterParams = z.infer<typeof TrackRegisterParamsSchema>;
......@@ -4,6 +4,7 @@ import { TeamPermission } from '../permission/user/controller';
import type { UserStatusEnum } from './constant';
import { TeamMemberStatusEnum } from './team/constant';
import { TeamTmbItemSchema } from './team/type';
import type { FastGPTSemType } from '../marketing/type';
export const UserTagsSchema = z.enum(['wecom']);
export const UserTagsEnum = UserTagsSchema.enum;
......@@ -26,9 +27,7 @@ export type UserModelSchema = {
status: `${UserStatusEnum}`;
lastLoginTmbId?: string;
passwordUpdateTime?: Date;
fastgpt_sem?: {
keyword: string;
};
fastgpt_sem?: FastGPTSemType;
contact?: string;
tags: UserTagsType[];
meta?: UserMetaType;
......
import { describe, expect, it } from 'vitest';
import {
DEFAULT_MAX_FOLDER_DEPTH,
canCreateFolderAtDepth,
canCreateSubFolder,
getCurrentFolderLevel,
isPathsSyncedWithParent,
normalizeParentId
} from '@fastgpt/global/common/parentFolder/depth';
describe('parent folder depth utils', () => {
describe('normalizeParentId', () => {
it('keeps non-empty string parent ids and treats other values as root', () => {
expect(normalizeParentId('parent-1')).toBe('parent-1');
expect(normalizeParentId('')).toBeNull();
expect(normalizeParentId(null)).toBeNull();
expect(normalizeParentId(undefined)).toBeNull();
expect(normalizeParentId(123)).toBeNull();
});
});
describe('isPathsSyncedWithParent', () => {
it('requires empty paths at root', () => {
expect(isPathsSyncedWithParent(null, [])).toBe(true);
expect(isPathsSyncedWithParent(undefined, [])).toBe(true);
expect(isPathsSyncedWithParent('', [])).toBe(true);
expect(isPathsSyncedWithParent(null, [{ parentId: 'stale' }])).toBe(false);
});
it('matches non-root parent id with the last path item', () => {
expect(
isPathsSyncedWithParent('parent-2', [{ parentId: 'parent-1' }, { parentId: 'parent-2' }])
).toBe(true);
expect(isPathsSyncedWithParent('parent-2', [])).toBe(false);
expect(isPathsSyncedWithParent('parent-2', [{ parentId: 'parent-1' }])).toBe(false);
expect(isPathsSyncedWithParent('parent-2', [{ parentId: null }])).toBe(false);
});
});
describe('getCurrentFolderLevel', () => {
it('returns 0 at root and paths length inside a folder', () => {
expect(getCurrentFolderLevel(null, 3)).toBe(0);
expect(getCurrentFolderLevel('', 3)).toBe(0);
expect(getCurrentFolderLevel('parent-1', 3)).toBe(3);
});
});
describe('canCreateFolderAtDepth', () => {
it('uses default max depth and handles the exact boundary', () => {
expect(DEFAULT_MAX_FOLDER_DEPTH).toBe(4);
expect(canCreateFolderAtDepth(0)).toBe(true);
expect(canCreateFolderAtDepth(3)).toBe(true);
expect(canCreateFolderAtDepth(4)).toBe(false);
});
it('respects a custom max folder level', () => {
expect(canCreateFolderAtDepth(1, 2)).toBe(true);
expect(canCreateFolderAtDepth(2, 2)).toBe(false);
});
});
describe('canCreateSubFolder', () => {
it('supports numeric paths length', () => {
expect(canCreateSubFolder(null, 10)).toBe(true);
expect(canCreateSubFolder('parent-3', 3)).toBe(true);
expect(canCreateSubFolder('parent-4', 4)).toBe(false);
expect(canCreateSubFolder('parent-2', 2, 2)).toBe(false);
});
it('allows creation when paths are stale and otherwise follows synced paths depth', () => {
expect(canCreateSubFolder('parent-2', [{ parentId: 'parent-1' }])).toBe(true);
expect(
canCreateSubFolder('parent-4', [
{ parentId: 'parent-1' },
{ parentId: 'parent-2' },
{ parentId: 'parent-3' },
{ parentId: 'parent-4' }
])
).toBe(false);
expect(
canCreateSubFolder('parent-3', [
{ parentId: 'parent-1' },
{ parentId: 'parent-2' },
{ parentId: 'parent-3' }
])
).toBe(true);
});
});
});
......@@ -8,7 +8,8 @@ import {
chatValue2RuntimePrompt,
runtimePrompt2ChatsValue,
getSystemPrompt_ChatItemType,
mergeAssistantFieldMessages
mergeAssistantFieldMessages,
normalizeAIChatValue
} from '@fastgpt/global/core/chat/adapt';
import { ChatRoleEnum, ChatFileTypeEnum } from '@fastgpt/global/core/chat/constants';
import { ChatCompletionRequestMessageRoleEnum } from '@fastgpt/global/core/ai/constants';
......@@ -792,6 +793,57 @@ describe('mergeAssistantFieldMessages', () => {
});
});
describe('normalizeAIChatValue', () => {
it('should merge plain text values and drop empty text placeholders when other values exist', () => {
expect(
normalizeAIChatValue([
{ text: { content: '' } },
{ text: { content: '' } },
{ reasoning: { content: 'thinking' }, text: { content: '' } },
{ text: { content: '' } },
{ text: { content: 'Final' } },
{ text: { content: ' answer' } },
{ text: { content: '' } },
{ text: { content: '' } }
])
).toEqual([
{ reasoning: { content: 'thinking' }, text: { content: '' } },
{ text: { content: 'Final answer' } }
]);
});
it('should keep one empty text value when all values are empty', () => {
expect(normalizeAIChatValue([])).toEqual([{ text: { content: '' } }]);
expect(normalizeAIChatValue([{ text: { content: '' } }, { text: { content: '' } }])).toEqual([
{ text: { content: '' } }
]);
});
it('should preserve reasoning boundaries for consecutive AI chat nodes', () => {
expect(
normalizeAIChatValue([
{
reasoning: { content: 'think 1' },
text: { content: 'answer 1' }
},
{
reasoning: { content: 'think 2' },
text: { content: 'answer 2' }
}
])
).toEqual([
{
reasoning: { content: 'think 1' },
text: { content: 'answer 1' }
},
{
reasoning: { content: 'think 2' },
text: { content: 'answer 2' }
}
]);
});
});
describe('chats2GPTMessages', () => {
it('should convert system message', () => {
const messages: ChatItemMiniType[] = [
......
......@@ -42,15 +42,23 @@ type CheckMoveFolderDepthProps = FolderDepthModelProps & {
isFolderType: FolderTypeChecker;
};
type DepthLimitOptions = {
maxAllowedDepth?: number;
limitErr?: CommonErrEnum;
};
/**
* 根据 parentId 向上追溯,计算父级目录深度。
* 根目录深度为 0;遇到 parentId 成环或父级不存在时拒绝请求。
* 传入 maxAllowedDepth 时,一旦已超过允许深度就直接抛错,避免继续无意义地向上扫描。
*/
const getParentFolderDepth = async ({
parentId,
teamId,
model
}: FolderDepthModelProps & { parentId: ParentIdType }): Promise<number> => {
model,
maxAllowedDepth,
limitErr = CommonErrEnum.invalidParams
}: FolderDepthModelProps & { parentId: ParentIdType } & DepthLimitOptions): Promise<number> => {
if (!parentId) return 0;
let depth = 0;
......@@ -71,6 +79,10 @@ const getParentFolderDepth = async ({
}
depth += 1;
if (maxAllowedDepth !== undefined && depth > maxAllowedDepth) {
throw limitErr;
}
currentId = doc.parentId ? String(doc.parentId) : null;
}
......@@ -80,16 +92,18 @@ const getParentFolderDepth = async ({
/**
* 计算被移动资源子树中文件夹的最大相对深度。
* 非文件夹资源返回 0;文件夹自身相对深度为 1。
* 传入 maxAllowedDepth 时,一旦子树相对深度超过目标剩余空间就直接抛错。
*/
const getSubtreeMaxFolderDepth = async ({
resourceId,
teamId,
model,
isFolderType
isFolderType,
maxAllowedDepth
}: FolderDepthModelProps & {
resourceId: string;
isFolderType: FolderTypeChecker;
}): Promise<number> => {
} & DepthLimitOptions): Promise<number> => {
const resource = await model.findById(resourceId, 'type teamId').lean<FolderResourceDoc>();
if (!resource || String(resource.teamId) !== String(teamId)) {
throw CommonErrEnum.invalidResource;
......@@ -108,6 +122,9 @@ const getSubtreeMaxFolderDepth = async ({
if (!current) break;
maxRelativeDepth = Math.max(maxRelativeDepth, current.depth);
if (maxAllowedDepth !== undefined && current.depth > maxAllowedDepth) {
throw CommonErrEnum.folderMoveDepthLimit;
}
const children = await model
.find({ parentId: current.id, teamId }, '_id type')
......@@ -143,7 +160,7 @@ const isInSubtree = async ({
while (currentId) {
if (currentId === ancestorId) return true;
if (visited.has(currentId)) return false;
if (visited.has(currentId)) throw CommonErrEnum.invalidParams;
visited.add(currentId);
const doc: FolderResourceDoc | null = await model
......@@ -166,11 +183,13 @@ export const checkCreateFolderDepth = async ({
model
}: CheckCreateFolderDepthProps) => {
const maxDepth = serviceEnv.MAX_FOLDER_DEPTH;
const parentDepth = await getParentFolderDepth({ parentId, teamId, model });
if (parentDepth + 1 > maxDepth) {
throw CommonErrEnum.folderDepthLimit;
}
await getParentFolderDepth({
parentId,
teamId,
model,
maxAllowedDepth: maxDepth - 1,
limitErr: CommonErrEnum.folderDepthLimit
});
};
/**
......@@ -190,6 +209,14 @@ export const checkMoveFolderDepth = async ({
throw CommonErrEnum.invalidParams;
}
const targetParentDepth = await getParentFolderDepth({
parentId: targetParentId,
teamId,
model,
maxAllowedDepth: maxDepth,
limitErr: CommonErrEnum.folderMoveDepthLimit
});
if (targetParentId) {
const movingIntoDescendant = await isInSubtree({
ancestorId: resourceId,
......@@ -202,12 +229,11 @@ export const checkMoveFolderDepth = async ({
}
}
const [targetParentDepth, subtreeMaxFolderDepth] = await Promise.all([
getParentFolderDepth({ parentId: targetParentId, teamId, model }),
getSubtreeMaxFolderDepth({ resourceId, teamId, model, isFolderType })
]);
if (targetParentDepth + subtreeMaxFolderDepth > maxDepth) {
throw CommonErrEnum.folderMoveDepthLimit;
}
await getSubtreeMaxFolderDepth({
resourceId,
teamId,
model,
isFolderType,
maxAllowedDepth: maxDepth - targetParentDepth
});
};
......@@ -23,6 +23,7 @@ export type AgentLoopChildrenInteractiveParams<TChildrenResponse = unknown> = {
export type AgentLoopToolChildrenInteractive<TChildrenResponse = unknown> = {
type: 'toolChildrenInteractive';
nodeResponseId?: string;
params: {
childrenResponse: TChildrenResponse;
toolParams: {
......
......@@ -51,26 +51,36 @@ export const ensureGenerateChat = async (params: EnsureGenerateChatParams) => {
);
};
/**
* 尝试占用一次会话生成槽。
*
* 同一个 `appId/chatId` 只允许一个请求进入生成中状态;已有 generating 记录时返回 false,
* 由 API 层转换为“当前会话正在运行”的错误。
*
* 这里用“条件匹配 + upsert + 唯一索引”实现无副作用抢占:只有非 generating
* 记录会被更新为 generating;如果已有 generating 记录,查询不会命中,upsert 会因
* `{ appId, chatId }` 唯一索引报 11000,再转换为 false。这样被拒绝的并发请求不会刷新
* updateTime 或覆盖 source/sourceName。
*/
export const tryStartGenerateChat = async (params: EnsureGenerateChatParams) => {
const { $set, $setOnInsert } = buildGeneratingChatUpdate(params);
try {
await MongoChat.updateOne(
await MongoChat.findOneAndUpdate(
{
appId: params.appId,
chatId: params.chatId,
chatGenerateStatus: {
$ne: ChatGenerateStatusEnum.generating
}
chatGenerateStatus: { $ne: ChatGenerateStatusEnum.generating }
},
{
$set,
$setOnInsert
},
{
upsert: true
upsert: true,
new: false
}
);
).lean();
return true;
} catch (error: any) {
......
......@@ -37,17 +37,6 @@ const ChatItemResponseSchema = new Schema({
// 按 chat item 拉取完整 nodeResponse rows;复合索引包含 _id,避免详情读取时额外排序。
ChatItemResponseSchema.index({ appId: 1, chatId: 1, chatItemDataId: 1, _id: 1 });
ChatItemResponseSchema.index(
{
appId: 1,
chatId: 1,
chatItemDataId: 1,
'data.id': 1
},
{
unique: true
}
);
// Clear expired response
ChatItemResponseSchema.index({ teamId: 1, time: -1 });
......
import type { ChatItemMiniType } from '@fastgpt/global/core/chat/type';
import { MongoChatItem } from './chatItemSchema';
import { MongoChat } from './chatSchema';
import { DispatchNodeResponseKeyEnum } from '@fastgpt/global/core/workflow/runtime/constants';
import { ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
import { MongoChatItemResponse } from './chatItemResponseSchema';
import type { ClientSession } from '../../common/mongo';
......@@ -12,14 +11,26 @@ import { composeChatItemResponseData } from './nodeResponseStorage';
const logger = getLogger(LogCategories.MODULE.CHAT.HISTORY);
/**
* 判断 AI 消息是否仍带旧的 chat_items.responseData 内联详情。
*
* 历史数据可能带空数组 responseData;只用 `.length` 会把这类旧记录误判成新记录,
* 进而用独立表的空结果覆盖原字段。新链路默认不写该字段,所以以字段是否存在作为边界。
*/
const hasInlineNodeResponses = (item: object) =>
Object.prototype.hasOwnProperty.call(item, DispatchNodeResponseKeyEnum.nodeResponse);
export type ChatItemNodeResponseMode = 'none' | 'preview' | 'full';
type ChatItemResponsePreviewProjection = Record<string, 1>;
const defaultNodeResponsePreviewProjection = {
chatItemDataId: 1,
'data.id': 1,
'data.parentId': 1,
'data.moduleType': 1,
'data.moduleName': 1,
'data.quoteList.id': 1,
'data.quoteList.collectionId': 1,
'data.quoteList.datasetId': 1,
'data.quoteList.sourceId': 1,
'data.quoteList.sourceName': 1,
'data.quoteList.chunkIndex': 1,
'data.quoteList.score': 1,
'data.toolId': 1,
'data.toolRes.citeLinks': 1,
'data.errorText': 1
} as const;
export async function getChatItems({
includeDeleted = false,
......@@ -27,6 +38,8 @@ export async function getChatItems({
chatId,
field,
limit,
nodeResponseMode,
nodeResponsePreviewProjection,
offset,
initialId,
......@@ -38,6 +51,8 @@ export async function getChatItems({
chatId?: string;
field: string;
limit: number;
nodeResponseMode?: ChatItemNodeResponseMode;
nodeResponsePreviewProjection?: ChatItemResponsePreviewProjection;
offset?: number;
initialId?: string;
......@@ -53,8 +68,9 @@ export async function getChatItems({
return { histories: [], total: 0, hasMorePrev: false, hasMoreNext: false };
}
// Extend dataId
const shouldReadNodeResponse = nodeResponseMode || 'none';
field = `dataId ${field}`;
const baseCondition = includeDeleted ? { appId, chatId } : { appId, chatId, deleteTime: null };
const { histories, total, hasMorePrev, hasMoreNext } = await (async () => {
......@@ -172,40 +188,51 @@ export async function getChatItems({
}
})();
// Add node responses field
if (field.includes(DispatchNodeResponseKeyEnum.nodeResponse) && histories.length > 0) {
if (shouldReadNodeResponse !== 'none' && histories.length > 0) {
const chatItemDataIds = histories
// 旧记录的 responseData 已内联在 chat_items,存在时不能被独立表的空结果覆盖。
.filter((item) => item.obj === ChatRoleEnum.AI && !hasInlineNodeResponses(item))
.filter((item) => item.obj === ChatRoleEnum.AI)
.map((item) => item.dataId);
if (chatItemDataIds.length > 0) {
const chatItemResponsesMap = await MongoChatItemResponse.find(
{ appId, chatId, chatItemDataId: { $in: chatItemDataIds } },
const isPreview = shouldReadNodeResponse === 'preview';
const rows = await MongoChatItemResponse.find(
{
appId,
chatId,
chatItemDataId: { $in: chatItemDataIds }
},
{
chatItemDataId: 1,
data: 1
...(isPreview
? nodeResponsePreviewProjection || defaultNodeResponsePreviewProjection
: { data: 1 })
}
)
.sort({ _id: 1 })
.lean()
.then((res) => {
const map = new Map<string, typeof res>();
res.forEach((item) => {
const val = map.get(item.chatItemDataId) || [];
val.push(item);
map.set(item.chatItemDataId, val);
});
return map;
.lean();
const chatItemResponsesMap = (() => {
const map = new Map<string, typeof rows>();
rows.forEach((item) => {
const val = map.get(item.chatItemDataId) || [];
val.push(item);
map.set(item.chatItemDataId, val);
});
return map;
})();
histories.forEach((item) => {
if (item.obj !== ChatRoleEnum.AI) return;
if (hasInlineNodeResponses(item)) return;
item.responseData = composeChatItemResponseData({
rows: chatItemResponsesMap.get(String(item.dataId)) || []
});
if (isPreview) {
item.responseData = chatItemResponsesMap
.get(String(item.dataId))
?.flatMap((row) => (row.data ? [row.data] : []));
} else {
item.responseData = composeChatItemResponseData({
rows: chatItemResponsesMap.get(String(item.dataId)) || []
});
}
});
}
}
......
import type { ChatItemMiniType, UserChatItemType } from '@fastgpt/global/core/chat/type';
import { UserError } from '@fastgpt/global/common/error/utils';
import { MongoChatItem } from './chatItemSchema';
import { MongoChatItem } from '../chatItemSchema';
import { ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
export const CHAT_DATA_ID_DUPLICATE_ERROR_MESSAGE = 'Chat dataId already exists';
/**
* 单轮对话进入工作流前的 dataId 校验参数。
*
* 当前运行前唯一性只约束 AI responseChatItemId;Human 消息允许与 AI 消息使用同一个
* dataId 表示同一轮对话。
*/
type ValidateChatRoundDataIdsParams = {
appId: string;
chatId: string;
......@@ -11,9 +18,11 @@ type ValidateChatRoundDataIdsParams = {
responseChatItemId?: string;
};
/** 过滤空值,避免未传 dataId 的旧调用或兼容数据参与重复判断。 */
const getValidDataIds = (dataIds: Array<string | undefined>) =>
dataIds.filter((dataId): dataId is string => typeof dataId === 'string' && dataId.length > 0);
/** 返回列表中第一个重复的 dataId,用于生成稳定、可读的错误信息。 */
const findDuplicateDataId = (dataIds: string[]) => {
const seen = new Set<string>();
......@@ -23,9 +32,15 @@ const findDuplicateDataId = (dataIds: string[]) => {
}
};
/** 从历史消息上下文中提取有效 dataId,供新请求进入前做重复检查。 */
export const getChatMessagesDataIds = (chatMessages: ChatItemMiniType[]) =>
getValidDataIds(chatMessages.map((item) => item.dataId));
/**
* 校验本次请求体内部不能携带重复 dataId。
*
* 这是纯内存检查,用于在访问数据库前快速拦截明显错误;旧数据中缺失 dataId 的消息会被忽略。
*/
export const assertNoDuplicateChatDataIdsInRequest = (dataIds: Array<string | undefined>) => {
const duplicateDataId = findDuplicateDataId(getValidDataIds(dataIds));
......@@ -34,6 +49,12 @@ export const assertNoDuplicateChatDataIdsInRequest = (dataIds: Array<string | un
}
};
/**
* 校验目标会话中是否已经存在任意相同 dataId 的 chat item。
*
* 这个方法不区分 Human/AI obj,适合通用历史消息场景;单轮工作流运行前的 AI response
* dataId 校验应使用 validateChatRoundDataIds。
*/
export const assertNoExistingChatDataIds = async ({
appId,
chatId,
......@@ -62,19 +83,32 @@ export const assertNoExistingChatDataIds = async ({
}
};
/**
* 校验本轮 AI responseChatItemId 是否已在当前会话中被占用。
*
* Human/AI 可以共用同一个 dataId 表示同一轮对话,所以这里仅检查 AI item,防止新的
* AI placeholder 或最终回复覆盖已有 AI 消息。
*/
export const validateChatRoundDataIds = async ({
appId,
chatId,
userContent,
responseChatItemId
}: ValidateChatRoundDataIdsParams) => {
const currentRoundDataIds = [userContent.dataId, responseChatItemId];
if (!responseChatItemId) return;
assertNoDuplicateChatDataIdsInRequest(currentRoundDataIds);
const existingChatItem = await MongoChatItem.findOne(
{
appId,
chatId,
obj: ChatRoleEnum.AI,
dataId: responseChatItemId
},
'dataId'
)
.lean()
.exec();
await assertNoExistingChatDataIds({
appId,
chatId,
dataIds: currentRoundDataIds
});
if (existingChatItem?.dataId) {
throw new UserError(`${CHAT_DATA_ID_DUPLICATE_ERROR_MESSAGE}: ${existingChatItem.dataId}`);
}
};
import { ChatErrEnum } from '@fastgpt/global/common/error/code/chat';
import { getNanoid } from '@fastgpt/global/common/string/tools';
import { ChatGenerateStatusEnum, ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
import type { ChatSourceEnum } from '@fastgpt/global/core/chat/constants';
import type { AIChatItemType, UserChatItemType } from '@fastgpt/global/core/chat/type';
import type { WorkflowInteractiveResponseType } from '@fastgpt/global/core/workflow/template/system/interactive/type';
import { mongoSessionRun } from '../../../common/mongo/sessionRun';
import { writePrimary } from '../../../common/mongo/utils';
import { MongoChatItem } from '../chatItemSchema';
import { MongoChat } from '../chatSchema';
import { tryStartGenerateChat, updateChatGenerateStatus } from '../chatGenerateStatus';
import { validateChatRoundDataIds } from './dataIdValidation';
import { getInteractiveResponseStatus } from '../interactiveResponseDataId';
export const NO_RECORD_CHAT_ID = 'NO_RECORD_HISTORIES';
/** 判断当前 chatId 是否为不落库的运行标记。 */
export const isSkipSaveChatId = (chatId?: string) => chatId === NO_RECORD_CHAT_ID;
const resolvePreChatRoundChatId = (chatId?: string) =>
chatId === NO_RECORD_CHAT_ID ? chatId : chatId || getNanoid(24);
/**
* 清理用户消息里的文件临时 URL,只保留 file key 参与持久化。
*
* 文件已经通过 key 持久化,URL 可能带 TTL 或签名信息,写入 chat item 会导致历史记录中保存
* 过期访问地址。
*/
export const stripUserContentFileUrls = (userContent: UserChatItemType & { dataId?: string }) => {
userContent.value.forEach((item) => {
if (item.file?.key) {
item.file.url = '';
}
});
};
export type EnsurePendingChatRoundParams = {
chatId: string;
appId: string;
teamId: string;
tmbId: string;
userContent: UserChatItemType & { dataId?: string };
responseChatItemId: string;
};
export type PrepareChatRoundParams = {
chatId: string;
appId: string;
teamId: string;
tmbId: string;
source: `${ChatSourceEnum}`;
sourceName?: string;
shareId?: string;
outLinkUid?: string;
userContent: UserChatItemType & { dataId?: string };
responseChatItemId: string;
};
export type PreChatRoundParams = Omit<PrepareChatRoundParams, 'chatId' | 'responseChatItemId'> & {
chatId?: string;
responseChatItemId?: string;
interactive?: WorkflowInteractiveResponseType;
};
export type PreChatRoundResult = {
chatId: string;
responseChatItemId: string;
shouldPersistChatRound: boolean;
shouldFinalizePreparedRound: boolean;
};
/**
* 读取预创建 Human/AI chat items 的 dataId。
*
* 新运行要求 prepare 阶段已经为本轮 Human/AI 创建同一个 dataId;缺失说明调用方绕过了
* preChatRound 或传入内容被错误覆盖,应直接失败,避免后续误写新记录。
*/
export const getPreparedRoundDataIds = ({
userContent,
aiContent
}: {
userContent: UserChatItemType & { dataId?: string };
aiContent: AIChatItemType & { dataId?: string };
}) => {
if (!userContent.dataId) {
throw new Error('Pending chat round human dataId is missing');
}
if (!aiContent.dataId) {
throw new Error('Pending chat round ai dataId is missing');
}
return {
humanDataId: userContent.dataId,
aiDataId: aiContent.dataId
};
};
/**
* 预创建一轮可保存的 Human/AI chat items。
*
* 这里使用严格 create,不再使用 upsert。调用方必须先确认 AI dataId 未被使用;
* Human 和 AI 使用同一个 roundDataId,便于客户端与服务端用一轮消息 ID 对齐。
*/
export const prepareChatRound = async (params: PrepareChatRoundParams) => {
const {
chatId,
appId,
teamId,
tmbId,
source,
sourceName,
shareId,
outLinkUid,
responseChatItemId
} = params;
if (isSkipSaveChatId(chatId)) return;
stripUserContentFileUrls(params.userContent);
params.userContent.dataId = responseChatItemId;
const now = new Date();
const userPayload: UserChatItemType & { dataId: string; obj: typeof ChatRoleEnum.Human } = {
...params.userContent,
dataId: responseChatItemId,
obj: ChatRoleEnum.Human
};
const aiPlaceholder: AIChatItemType & { dataId: string } = {
dataId: responseChatItemId,
obj: ChatRoleEnum.AI,
value: []
};
await mongoSessionRun(async (session) => {
await MongoChat.updateOne(
{
appId,
chatId
},
{
$set: {
teamId,
tmbId,
appId,
chatId,
source,
sourceName,
shareId,
outLinkUid,
updateTime: now,
hasBeenRead: false,
chatGenerateStatus: ChatGenerateStatusEnum.generating
},
$setOnInsert: {
createTime: now
}
},
{
session,
upsert: true
}
);
await MongoChatItem.create(
[
{
teamId,
tmbId,
chatId,
appId,
...userPayload
},
{
teamId,
tmbId,
chatId,
appId,
...aiPlaceholder
}
],
{ session, ordered: true, ...writePrimary }
);
});
};
/**
* 业务入口进入 workflow 前的唯一准备方法。
*
* 它负责解析最终 chatId/responseChatItemId、占用生成槽、检查 AI dataId 冲突,并在
* 需要持久化时预创建本轮 Human/AI placeholder。失败时如果已经占用生成槽,会立刻将
* chatGenerateStatus 标记为 error,避免会话长期停留在 generating。
*/
export const preChatRound = async (params: PreChatRoundParams): Promise<PreChatRoundResult> => {
const chatId = resolvePreChatRoundChatId(params.chatId);
const responseChatItemId = params.responseChatItemId || getNanoid(24);
const shouldPersistChatRound = !isSkipSaveChatId(chatId);
const interactiveStatus = getInteractiveResponseStatus({
interactive: params.interactive,
userContent: params.userContent
});
const isInteractiveContinue = !!params.interactive && interactiveStatus !== 'query';
if (!shouldPersistChatRound) {
return {
chatId,
responseChatItemId,
shouldPersistChatRound: false,
shouldFinalizePreparedRound: false
};
}
const canStartGenerate = await tryStartGenerateChat({
appId: params.appId,
chatId,
teamId: params.teamId,
tmbId: params.tmbId,
source: params.source,
sourceName: params.sourceName,
shareId: params.shareId,
outLinkUid: params.outLinkUid
});
if (!canStartGenerate) {
throw ChatErrEnum.chatIsGenerating;
}
try {
if (isInteractiveContinue) {
const previousAiItem = await MongoChatItem.findOne(
{
appId: params.appId,
chatId,
obj: ChatRoleEnum.AI
},
'dataId'
)
.sort({ _id: -1 })
.lean()
.exec();
if (!previousAiItem?.dataId) {
throw new Error(`Interactive continue chat item not found: ${chatId}`);
}
return {
chatId,
responseChatItemId: previousAiItem.dataId,
shouldPersistChatRound: true,
shouldFinalizePreparedRound: false
};
}
await validateChatRoundDataIds({
appId: params.appId,
chatId,
userContent: params.userContent,
responseChatItemId
});
await prepareChatRound({
...params,
chatId,
responseChatItemId
});
return {
chatId,
responseChatItemId,
shouldPersistChatRound: true,
shouldFinalizePreparedRound: true
};
} catch (error) {
await updateChatGenerateStatus({
appId: params.appId,
chatId,
status: ChatGenerateStatusEnum.error
});
throw error;
}
};
......@@ -134,8 +134,7 @@ export const dispatchLoop = async (props: Props): Promise<Response> => {
[DispatchNodeResponseKeyEnum.nodeResponse]: {
totalPoints,
loopInput: loopInputArray,
loopResult: outputValueArr,
mergeSignId: props.node.nodeId
loopResult: outputValueArr
},
[DispatchNodeResponseKeyEnum.customFeedbacks]:
customFeedbacks.length > 0 ? customFeedbacks : undefined
......
import { getNanoid } from '@fastgpt/global/common/string/tools';
import type { ChatHistoryItemResType } from '@fastgpt/global/core/chat/type';
import { stripChildTotalPoints } from '@fastgpt/global/core/chat/utils';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import type { ChatNodeUsageType } from '@fastgpt/global/support/wallet/bill/type';
import type { AgentLoopEvent, AgentLoopToolCatalog } from '../../../../../ai/llm/agentLoop';
......@@ -92,9 +91,6 @@ export const useToolNodeResponse = ({
};
};
const stripNodeResponseChildTotalPoints = (nodeResponse: ChatHistoryItemResType) =>
stripChildTotalPoints(nodeResponse);
const getUpdatePlanStatus = (call: ToolResponseEvent['call']): AgentPlanStatus => {
/**
* update_plan 既可能是首次设置/替换计划,也可能是状态更新。
......@@ -180,19 +176,17 @@ export const useToolNodeResponse = ({
: undefined;
const childrenResponses = compressNodeResponse ? [compressNodeResponse] : [];
appendNodeResponse(
stripNodeResponseChildTotalPoints({
id: `${node.nodeId}-plan-${call.id}`,
nodeId: `${node.nodeId}-plan-${call.id}`,
moduleName: AgentNodeResponseDisplay.plan.moduleName,
moduleType: node.flowNodeType,
moduleLogo: AgentNodeResponseDisplay.plan.moduleLogo,
runningTime: seconds,
textOutput: response,
agentPlanStatus,
...(childrenResponses.length > 0 ? { childrenResponses } : {})
})
);
appendNodeResponse({
id: `${node.nodeId}-plan-${call.id}`,
nodeId: `${node.nodeId}-plan-${call.id}`,
moduleName: AgentNodeResponseDisplay.plan.moduleName,
moduleType: node.flowNodeType,
moduleLogo: AgentNodeResponseDisplay.plan.moduleLogo,
runningTime: seconds,
textOutput: response,
agentPlanStatus,
...(childrenResponses.length > 0 ? { childrenResponses } : {})
});
};
const createToolNodeResponse = (event: ToolResponseEvent): ChatHistoryItemResType => {
......@@ -222,12 +216,12 @@ export const useToolNodeResponse = ({
...(compressNodeResponse ? [compressNodeResponse] : [])
];
return stripNodeResponseChildTotalPoints({
return {
...toolNodeResponse,
runningTime: toolNodeResponse.runningTime ?? event.seconds,
toolRes: toolNodeResponse.toolRes ?? event.response,
...(childrenResponses.length > 0 ? { childrenResponses } : {})
});
};
};
/**
......
......@@ -158,7 +158,6 @@ export const dispatchRunTools = async (props: DispatchToolModuleProps): Promise<
model: modelName,
query: userChatInput,
historyPreview,
mergeSignId: nodeId,
finishReason: finish_reason,
llmRequestIds: requestIds
};
......
......@@ -224,8 +224,7 @@ export const dispatchRunAppNode = async (props: Props): Promise<Response> => {
totalPoints: usagePoints,
query: userChatInput,
textOutput: text,
childResponseCount,
mergeSignId: props.node.nodeId
childResponseCount
},
[DispatchNodeResponseKeyEnum.toolResponse]: text,
[DispatchNodeResponseKeyEnum.customFeedbacks]: customFeedbacks
......
......@@ -16,6 +16,7 @@ import {
DispatchNodeResponseKeyEnum,
SseResponseEventEnum
} from '@fastgpt/global/core/workflow/runtime/constants';
import { normalizeAIChatValue } from '@fastgpt/global/core/chat/adapt';
import type {
ChatDispatchProps,
DispatchNodeResultType,
......@@ -23,7 +24,7 @@ import type {
} from '@fastgpt/global/core/workflow/runtime/type';
import type { RuntimeNodeItemType } from '@fastgpt/global/core/workflow/runtime/type';
import { getErrText, UserError } from '@fastgpt/global/common/error/utils';
import { filterNodeResponseTreeData, stripChildTotalPoints } from '@fastgpt/global/core/chat/utils';
import { filterNodeResponseTreeData } from '@fastgpt/global/core/chat/utils';
import {
filterWorkflowEdges,
textAdaptGptResponse,
......@@ -94,7 +95,7 @@ type Props = Omit<
| 'variableState'
| 'responseChatItemId'
> & {
responseChatItemId?: string;
responseChatItemId: string;
variables: Record<string, any>;
runtimeNodes: RuntimeNodeItemType[];
runtimeEdges: RuntimeEdgeItemType[];
......@@ -137,7 +138,7 @@ export async function dispatchWorkFlow({
chatId,
apiVersion
} = data;
const responseChatItemId = data.responseChatItemId || getNanoid(24);
const responseChatItemId = data.responseChatItemId;
// Check url valid
const invalidInput = query.some((item) => {
......@@ -259,7 +260,6 @@ export async function dispatchWorkFlow({
: undefined;
const { nodeResponseWriter } = await createWorkflowEntryNodeResponseWriter({
lastInteractive: data.lastInteractive,
teamId: data.runningAppInfo.teamId,
appId: data.runningAppInfo.id,
chatId,
......@@ -369,6 +369,7 @@ export class WorkflowQueue {
| {
entryNodeIds: string[];
interactiveResponse: InteractiveNodeResponseType;
nodeResponseId?: string;
}
| undefined;
system_memories: Record<string, any> = {}; // Workflow node memories
......@@ -800,6 +801,7 @@ export class WorkflowQueue {
private async nodeRunWithActive(node: RuntimeNodeItemType): Promise<{
node: RuntimeNodeItemType;
runStatus: 'run';
nodeResponseId: string;
result: NodeResponseCompleteType;
}> {
const mode = this.isDebugMode ? 'test' : this.data.mode;
......@@ -809,8 +811,11 @@ export class WorkflowQueue {
};
const executeNode = async (stepSpan?: Span): Promise<WorkflowObservedStepResult> => {
const nodeResponseId = getNanoid();
const nodeResponseId =
this.data.lastInteractive?.nodeResponseId &&
this.data.lastInteractive.entryNodeIds?.includes(node.nodeId)
? this.data.lastInteractive.nodeResponseId
: getNanoid();
// push run status messages
if (node.showStatus && !this.data.isToolCall) {
this.data.workflowStreamResponse?.({
......@@ -987,7 +992,7 @@ export class WorkflowQueue {
filteredResponses.forEach((item) => {
this.data.workflowStreamResponse?.({
event: SseResponseEventEnum.flowNodeResponse,
data: stripChildTotalPoints(item)
data: item
});
});
}
......@@ -1025,6 +1030,7 @@ export class WorkflowQueue {
return {
node,
runStatus: 'run',
nodeResponseId,
result: {
...dispatchRes,
runtimeNodeResponseSummary: mergeRuntimeNodeResponseSummary(
......@@ -1380,6 +1386,10 @@ export class WorkflowQueue {
if (this.isDebugMode) {
this.debugNextStepRunNodes = this.debugNextStepRunNodes.concat([nodeRunResult.node]);
}
const nodeResponseId =
nodeRunResult.runStatus === 'run'
? nodeRunResult.nodeResponseId
: nodeRunResult.result[DispatchNodeResponseKeyEnum.nodeResponse]?.id;
// For the pause interactive response, there may be multiple nodes triggered at the same time, so multiple entry nodes need to be recorded.
// For other interactive nodes, only one will be triggered at the same time.
......@@ -1388,12 +1398,14 @@ export class WorkflowQueue {
entryNodeIds: this.nodeInteractiveResponse?.entryNodeIds
? this.nodeInteractiveResponse.entryNodeIds.concat(nodeRunResult.node.nodeId)
: [nodeRunResult.node.nodeId],
interactiveResponse
interactiveResponse,
nodeResponseId
};
} else {
this.nodeInteractiveResponse = {
entryNodeIds: [nodeRunResult.node.nodeId],
interactiveResponse
interactiveResponse,
nodeResponseId
};
}
return;
......@@ -1410,10 +1422,12 @@ export class WorkflowQueue {
/* Have interactive result, computed edges and node outputs */
handleInteractiveResult({
entryNodeIds,
interactiveResponse
interactiveResponse,
nodeResponseId
}: {
entryNodeIds: string[];
interactiveResponse: InteractiveNodeResponseType;
nodeResponseId?: string;
}): AIChatItemValueItemType {
// Get node outputs
const nodeOutputs: NodeOutputItemType[] = [];
......@@ -1431,6 +1445,7 @@ export class WorkflowQueue {
const interactiveResult: WorkflowInteractiveResponseType = {
...interactiveResponse,
...(nodeResponseId ? { nodeResponseId } : {}),
skipNodeQueue: Array.from(this.skipNodeQueue.values()).map((item) => ({
id: item.node.nodeId,
skippedNodeIdList: Array.from(item.skippedNodeIdList)
......@@ -1477,7 +1492,6 @@ export class WorkflowQueue {
}
}
export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowResponse> => {
data.responseChatItemId ||= getNanoid(24);
// Over max depth
const previousWorkflowDispatchDeep = data.workflowDispatchDeep;
const currentWorkflowDispatchDeep = previousWorkflowDispatchDeep + 1;
......@@ -1582,7 +1596,8 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR
if (workflowQueue.nodeInteractiveResponse) {
const interactiveAssistant = workflowQueue.handleInteractiveResult({
entryNodeIds: workflowQueue.nodeInteractiveResponse.entryNodeIds,
interactiveResponse: workflowQueue.nodeInteractiveResponse.interactiveResponse
interactiveResponse: workflowQueue.nodeInteractiveResponse.interactiveResponse,
nodeResponseId: workflowQueue.nodeInteractiveResponse.nodeResponseId
});
if (workflowQueue.isRootRuntime) {
workflowQueue.chatAssistantResponse.push(interactiveAssistant);
......@@ -1613,7 +1628,7 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR
debugResponse: workflowQueue.getDebugResponse(),
workflowInteractiveResponse: interactiveResult,
[DispatchNodeResponseKeyEnum.runTimes]: workflowQueue.workflowRunTimes,
[DispatchNodeResponseKeyEnum.assistantResponses]: mergeAssistantResponseAnswerText(
[DispatchNodeResponseKeyEnum.assistantResponses]: normalizeAIChatValue(
workflowQueue.chatAssistantResponse
),
[DispatchNodeResponseKeyEnum.toolResponse]: workflowQueue.toolRunResponse,
......@@ -1644,54 +1659,3 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR
}
});
};
/**
* 合并连续的纯文本 assistant value,减少普通 answer 片段在页面上分裂展示。
*
* 带 reasoning、工具、交互、计划等结构化信息的 value 必须保留独立边界;否则连续
* AI 节点的第二段 reasoning 会在合并 text 时被吞掉,刷新历史后页面看不到对应思考。
*/
export const mergeAssistantResponseAnswerText = (response: AIChatItemValueItemType[]) => {
const result: AIChatItemValueItemType[] = [];
const isPlainTextValue = (item: AIChatItemValueItemType) =>
!!item.text &&
!item.id &&
!item.planId &&
!item.reasoning &&
!item.tools &&
!item.skills &&
!item.interactive &&
!item.plan &&
!item.planStatus &&
!item.agentPlanUpdate &&
!item.agentAsk &&
!item.agentStopGate &&
!item.contextCheckpoint &&
!item.tool &&
!item.hideReason &&
!item.hideInUI;
// 合并连续的text
for (let i = 0; i < response.length; i++) {
const item = response[i];
if (isPlainTextValue(item)) {
const text = item.text?.content || '';
const lastItem = result[result.length - 1];
if (lastItem && isPlainTextValue(lastItem) && lastItem.text?.content) {
lastItem.text.content += text;
continue;
}
}
result.push(item);
}
// If result is empty, auto add a text message
if (result.length === 0) {
result.push({
text: { content: '' }
});
}
return result;
};
......@@ -78,7 +78,6 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
return getNodeErrResponse({
error: preCheckError,
responseData: {
mergeSignId: node.nodeId,
...(mode === LoopRunModeEnum.array ? { loopRunInput: inputArray } : {})
}
});
......@@ -113,13 +112,24 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
// Pre-interrupt runtime summary survives across resume here, so loopRun can still
// calculate finished nodes, child stats and wall time without retaining full child details.
let pendingIterationSummary = interactiveData?.pendingIterationSummary;
const getIncrementalChildResponseCount = (
summary: ReturnType<typeof getRuntimeNodeResponseSummary>,
preInterruptSummary: typeof pendingIterationSummary
) => {
if (!preInterruptSummary) return summary.childResponseCount;
const getWrapperSummary = ({
isResumeIteration,
response
}: {
isResumeIteration: boolean;
response: Response;
}) => {
const currentSummary = getRuntimeNodeResponseSummary(response);
const fullSummary = isResumeIteration
? mergeRuntimeNodeResponseSummary(pendingIterationSummary, currentSummary)
: currentSummary;
return (summary.childResponseCount || 0) - (preInterruptSummary.childResponseCount || 0);
// 同一个 iterationResponseId 会在暂停和恢复后各写一条 wrapper row;读取时数值字段按
// id/parentId 累加,所以恢复后的 wrapper 必须只写本次 resume 产生的增量统计。
return {
fullSummary,
wrapperSummary: isResumeIteration ? currentSummary : fullSummary
};
};
const resumeIteration = interactiveData?.iteration;
......@@ -152,7 +162,8 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
}
const isResumeIteration = !!interactiveData && iteration === resumeIteration;
const iterationResponseId = `${node.nodeId}_iter_${iteration}`;
const loopRunNodeResponseId = props.nodeResponseParentId || node.nodeId;
const iterationResponseId = `${loopRunNodeResponseId}:iter:${iteration}`;
if (isResumeIteration) {
isolatedNodes.forEach((n) => {
......@@ -183,17 +194,12 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
// Merge pre-interrupt runtime summary so resumed iteration still sees the full
// set of finished nodes and stats without keeping full child nodeResponse data.
const iterationSummary = isResumeIteration
? mergeRuntimeNodeResponseSummary(
pendingIterationSummary,
getRuntimeNodeResponseSummary(response)
)
: getRuntimeNodeResponseSummary(response);
const iterationChildResponseCount = getIncrementalChildResponseCount(
iterationSummary,
isResumeIteration ? pendingIterationSummary : undefined
);
const iterationRunningTime = iterationSummary.runningTime;
const { fullSummary: iterationSummary, wrapperSummary } = getWrapperSummary({
isResumeIteration,
response
});
const iterationChildResponseCount = wrapperSummary.childResponseCount;
const iterationRunningTime = wrapperSummary.runningTime;
assistantResponses.push(...response.assistantResponses);
const iterationTotalPoints = pushSubWorkflowUsage({
usagePush: props.usagePush,
......@@ -201,13 +207,12 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
name,
iteration
});
const iterationDetailTotalPoints = iterationSummary.totalPoints ?? iterationTotalPoints;
const iterationDetailTotalPoints = wrapperSummary.totalPoints ?? iterationTotalPoints;
totalPoints += iterationTotalPoints;
collectResponseFeedbacks(response, customFeedbacks);
// Apply `finishedNodeIds` over the merged children so pre-interrupt nodes count
// as finished for the customOutputs snapshot. 交互暂停时也要先算一次,避免
// partial iteration wrapper 缺少已完成输出。
// as finished for the customOutputs snapshot.
const finishedNodeIds = new Set(iterationSummary.finishedNodeIds);
const customOutputs = readCustomOutputSnapshot({
customOutputInputs,
......@@ -241,7 +246,9 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
};
// Pause: stash accumulated children so the next resume still sees pre-interrupt
// nodes (supports multiple interrupts in the same iteration).
// nodes (supports multiple interrupts in the same iteration). 暂停时也写一次 wrapper,
// 作为本轮 child nodeResponses 的 parent;恢复完成后会用同一个 iterationResponseId
// 再写增量 wrapper,读取时按 id/parentId 累加统计并合并 children。
if (response.workflowInteractiveResponse) {
interactiveResponse = response.workflowInteractiveResponse;
pendingIterationSummary = iterationSummary;
......@@ -315,7 +322,6 @@ export const dispatchLoopRun = async (props: Props): Promise<Response> => {
loopRunIterations: loopHistory.length,
loopRunHistory: loopHistory,
childResponseCount,
mergeSignId: node.nodeId,
...(errorText ? { errorText } : {})
},
[DispatchNodeResponseKeyEnum.customFeedbacks]:
......
......@@ -192,8 +192,7 @@ export const dispatchParallelRun = async (props: Props): Promise<Response> => {
parallelInput: loopInputArray,
parallelResult: filteredArray,
parallelRunDetail: fullDetail,
childResponseCount: rootChildResponseCount,
mergeSignId: node.nodeId
childResponseCount: rootChildResponseCount
},
[DispatchNodeResponseKeyEnum.customFeedbacks]:
customFeedbacks.length > 0 ? customFeedbacks : undefined
......
import type { ChatDispatchProps } from '@fastgpt/global/core/workflow/runtime/type';
import type { WorkflowNodeResponseWriter } from '../../../chat/nodeResponseStorage';
import { createWorkflowNodeResponseWriter } from '../../../chat/nodeResponseStorage';
......@@ -16,14 +15,12 @@ export type WorkflowNodeResponseWriteConfig = {
* 推断。子 workflow 只复用这个 writer,不关心当前请求到底落库还是仅保留请求内 flat 数据。
*/
export const createWorkflowEntryNodeResponseWriter = async ({
lastInteractive,
teamId,
appId,
chatId,
chatItemDataId,
nodeResponseWriteConfig
}: {
lastInteractive?: ChatDispatchProps['lastInteractive'];
teamId: string;
appId: string;
chatId: string;
......@@ -34,7 +31,6 @@ export const createWorkflowEntryNodeResponseWriter = async ({
}> => {
return {
nodeResponseWriter: await createWorkflowNodeResponseWriter({
mode: lastInteractive ? 'append' : 'replace',
teamId,
appId,
chatId,
......
import { getErrText } from '@fastgpt/global/common/error/utils';
import { ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
import type { ChatHistoryItemResType, ChatItemMiniType } from '@fastgpt/global/core/chat/type';
import { getChildrenResponses, hasContextCheckpoint } from '@fastgpt/global/core/chat/utils';
import { hasContextCheckpoint } from '@fastgpt/global/core/chat/utils';
import { getChildrenResponses } from '@fastgpt/global/core/chat/utils/mergeNode';
import type { ChatNodeUsageType } from '@fastgpt/global/support/wallet/bill/type';
import type { DispatchFlowResponse, RuntimeNodeResponseSummary } from '../type';
import { NodeInputKeyEnum, NodeOutputKeyEnum } from '@fastgpt/global/core/workflow/constants';
......
......@@ -7,6 +7,7 @@ import type { RuntimeNodeItemType } from '@fastgpt/global/core/workflow/runtime/
export type WorkflowObservedStepResult = {
node: RuntimeNodeItemType;
runStatus: 'run';
nodeResponseId: string;
result: {
[DispatchNodeResponseKeyEnum.nodeResponse]?: ChatHistoryItemResType;
error?: {
......
import { ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
import type { UserChatItemValueItemType } from '@fastgpt/global/core/chat/type';
import { ChatGenerateStatusEnum, ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
import type { UserChatItemType, UserChatItemValueItemType } from '@fastgpt/global/core/chat/type';
import {
getWorkflowEntryNodeIds,
getMaxHistoryLimitFromNodes,
......@@ -10,7 +10,13 @@ import type { OutlinkAppType, OutLinkSchemaType } from '@fastgpt/global/support/
import { getAppLatestVersion } from '../../../core/app/version/controller';
import { MongoApp } from '../../../core/app/schema';
import { getChatItems } from '../../../core/chat/controller';
import { pushChatRecords } from '../../../core/chat/saveChat';
import {
failChatRound,
finalizeChatRound,
type Props as SaveChatProps
} from '../../../core/chat/saveChat';
import { preChatRound, type PreChatRoundResult } from '../../../core/chat/utils/prepare';
import { updateChatGenerateStatus } from '../../../core/chat/chatGenerateStatus';
import { dispatchWorkFlow } from '../../../core/workflow/dispatch';
import { getRunningUserInfoByTmbId } from '../../../support/user/team/utils';
import { SseResponseEventEnum } from '@fastgpt/global/core/workflow/runtime/constants';
......@@ -100,6 +106,11 @@ export async function outlinkInvokeChat<T extends OutlinkAppType>({
streamId
}: outLinkInvokeChatProps<T>) {
const streamResKey = `${STREAM_CACHE_KEY_PREFIX}${streamId}`;
const roundState = {
preparedRound: undefined as PreChatRoundResult | undefined,
appId: '',
finalized: false
};
try {
// Get app workflow config
......@@ -190,7 +201,25 @@ export async function outlinkInvokeChat<T extends OutlinkAppType>({
// Merge global variables from database
const variables = chatDetail?.variables ?? {};
const responseChatItemId = getNanoid(24);
const userContent: UserChatItemType & { dataId?: string } = {
dataId: messageId,
obj: ChatRoleEnum.Human,
value: query
};
const preparedRound = await preChatRound({
appId: String(app._id),
chatId,
teamId: String(outLinkConfig.teamId),
tmbId: String(outLinkConfig.tmbId),
source: getChatSourceByPublishChannel(outLinkConfig.type),
sourceName: outLinkConfig.name,
shareId: outLinkConfig.shareId,
outLinkUid: chatUserId,
userContent,
responseChatItemId: messageId
});
roundState.preparedRound = preparedRound;
roundState.appId = String(app._id);
const {
assistantResponses,
......@@ -212,8 +241,8 @@ export async function outlinkInvokeChat<T extends OutlinkAppType>({
},
runningUserInfo: await getRunningUserInfoByTmbId(app.tmbId),
uid: chatUserId || outLinkConfig.tmbId,
chatId,
responseChatItemId,
chatId: preparedRound.chatId,
responseChatItemId: preparedRound.responseChatItemId,
variables,
histories,
query: query,
......@@ -261,9 +290,9 @@ export async function outlinkInvokeChat<T extends OutlinkAppType>({
})();
// Save and reply
await pushChatRecords({
chatId,
appId: app._id,
const saveParams: SaveChatProps = {
chatId: preparedRound.chatId,
appId: String(app._id),
teamId: outLinkConfig.teamId,
tmbId: outLinkConfig.tmbId,
outLinkUid: chatUserId,
......@@ -274,13 +303,9 @@ export async function outlinkInvokeChat<T extends OutlinkAppType>({
shareId: outLinkConfig.shareId,
source: getChatSourceByPublishChannel(outLinkConfig.type),
sourceName: outLinkConfig.name,
userContent: {
dataId: messageId,
obj: ChatRoleEnum.Human,
value: query
},
userContent,
aiContent: {
dataId: responseChatItemId,
dataId: preparedRound.responseChatItemId,
obj: ChatRoleEnum.AI,
value: assistantResponses,
memories: system_memories
......@@ -289,7 +314,9 @@ export async function outlinkInvokeChat<T extends OutlinkAppType>({
durationSeconds,
errorMsg: replyResult?.success ? undefined : replyResult?.errmsg,
nodeResponseSummary
});
};
await finalizeChatRound(saveParams);
roundState.finalized = true;
const totalPoints = flowUsages.reduce((sum, item) => sum + (item.totalPoints || 0), 0);
addOutLinkUsage({
......@@ -301,6 +328,38 @@ export async function outlinkInvokeChat<T extends OutlinkAppType>({
await appendRedisCache(streamResKey, STREAM_END_FLAG, 60);
}
} catch (error) {
const { preparedRound } = roundState;
if (!roundState.finalized && preparedRound?.shouldPersistChatRound && roundState.appId) {
if (preparedRound.shouldFinalizePreparedRound) {
await failChatRound({
appId: roundState.appId,
chatId: preparedRound.chatId,
responseChatItemId: preparedRound.responseChatItemId,
error
}).catch((saveError) => {
logger.error('Outlink invoke chat mark error failed', {
shareId: outLinkConfig.shareId,
chatId,
messageId,
error: saveError
});
});
} else {
await updateChatGenerateStatus({
appId: roundState.appId,
chatId: preparedRound.chatId,
status: ChatGenerateStatusEnum.error
}).catch((saveError) => {
logger.error('Outlink invoke chat unlock failed', {
shareId: outLinkConfig.shareId,
chatId,
messageId,
error: saveError
});
});
}
}
logger.error('Outlink invoke chat failed', {
shareId: outLinkConfig.shareId,
chatId,
......
......@@ -62,7 +62,6 @@ const UserSchema = new Schema({
ref: userCollectionName
},
fastgpt_sem: Object,
sourceDomain: String,
phonePrefix: Number,
contact: String,
......
......@@ -323,6 +323,71 @@ describe('runAgentLoop with mocked createLLMResponse', () => {
});
});
it('keeps child nodeResponseId only inside childrenResponse when a tool call pauses', async () => {
const onRunTool = vi.fn(async () => ({
response: 'waiting for user selection',
assistantMessages: [],
usages: [],
interactive: {
type: 'userSelect',
nodeResponseId: 'child_select_response',
entryNodeIds: ['select_node'],
memoryEdges: [],
nodeOutputs: [],
params: {
description: 'Choose one',
userSelectOptions: []
}
}
}));
mockCreateLLMResponseQueue(createLLMResponseMock, [
toolCall({
id: 'call_search',
name: 'search',
args: {
q: 'FastGPT'
}
}),
text({ requestId: 'req_should_not_be_used', content: 'unused' })
]);
const result = await runAgentLoop({
maxRunAgentTimes: 5,
body: {
model: 'gpt-4',
stream: true,
messages: [
{
role: ChatCompletionRequestMessageRoleEnum.User,
content: 'search FastGPT'
}
],
tools: [searchTool]
},
usagePush: vi.fn(),
isAborted: () => false,
onRunTool,
onRunInteractiveTool: vi.fn()
});
expect(createLLMResponseMock).toHaveBeenCalledTimes(1);
expect(result.interactiveResponse).toEqual(
expect.objectContaining({
type: 'toolChildrenInteractive',
params: expect.objectContaining({
childrenResponse: expect.objectContaining({
nodeResponseId: 'child_select_response'
}),
toolParams: expect.objectContaining({
toolCallId: 'call_search'
})
})
})
);
expect(result.interactiveResponse?.nodeResponseId).toBeUndefined();
});
it('feeds the compressed tool response into the next LLM request', async () => {
compressToolResponseMock.mockImplementation(async ({ response }) => ({
compressed: `compressed:${response}`
......
......@@ -768,7 +768,8 @@ describe('getChatItems', () => {
chatId,
offset: 0,
limit: 10,
field: 'obj value responseData'
field: 'obj value',
nodeResponseMode: 'full'
});
expect(result.histories).toHaveLength(1);
......@@ -784,7 +785,7 @@ describe('getChatItems', () => {
});
});
it('keeps inline chatItem responseData for legacy records', async () => {
it('reads persisted rows even when chat item still has legacy inline responseData', async () => {
await MongoChatItem.create({
teamId: testUser.teamId,
tmbId: testUser.tmbId,
......@@ -822,13 +823,14 @@ describe('getChatItems', () => {
chatId,
offset: 0,
limit: 10,
field: 'obj value responseData'
field: 'obj value',
nodeResponseMode: 'full'
});
expect(result.histories[0].responseData?.map((item) => item.id)).toEqual(['fallback-root']);
expect(result.histories[0].responseData?.map((item) => item.id)).toEqual(['persisted-root']);
});
it('does not overwrite legacy chatItem responseData even when it is empty', async () => {
it('does not use empty legacy chat item responseData as a fallback', async () => {
await MongoChatItem.create({
teamId: testUser.teamId,
tmbId: testUser.tmbId,
......@@ -859,10 +861,151 @@ describe('getChatItems', () => {
chatId,
offset: 0,
limit: 10,
field: 'obj value responseData'
field: 'obj value',
nodeResponseMode: 'full'
});
expect(result.histories[0].responseData).toEqual([]);
expect(result.histories[0].responseData?.map((item) => item.id)).toEqual(['persisted-root']);
});
it('loads lightweight preview response rows without composing responseData tree', async () => {
const citedQuoteId = '0123456789abcdef01234567';
const uncitedQuoteId = 'fedcba9876543210fedcba98';
const aiDataId = 'preview-ai-data-id';
await MongoChatItem.create({
teamId: testUser.teamId,
tmbId: testUser.tmbId,
userId: testUser.userId,
appId,
chatId,
dataId: aiDataId,
obj: ChatRoleEnum.AI,
value: [
{
type: 'text',
text: {
content: `Answer with cite [${citedQuoteId}](CITE)`
}
}
],
responseData: [
{
id: 'legacy-root',
moduleName: 'Legacy',
moduleType: FlowNodeTypeEnum.chatNode,
errorText: 'legacy error'
}
]
});
await MongoChatItemResponse.create([
{
teamId: testUser.teamId,
appId,
chatId,
chatItemDataId: aiDataId,
data: {
id: 'dataset-response',
moduleName: 'Dataset',
moduleType: FlowNodeTypeEnum.datasetSearchNode,
quoteList: [
{
id: citedQuoteId,
datasetId: 'dataset-1',
collectionId: 'collection-1',
sourceId: 'source-1',
sourceName: 'source.md',
chunkIndex: 0,
score: 0.9
},
{
id: uncitedQuoteId,
datasetId: 'dataset-1',
collectionId: 'collection-1',
sourceId: 'source-2',
sourceName: 'unused.md',
chunkIndex: 1,
score: 0.5
}
]
}
},
{
teamId: testUser.teamId,
appId,
chatId,
chatItemDataId: aiDataId,
data: {
id: 'tool-response',
moduleName: 'Tool',
moduleType: FlowNodeTypeEnum.tool,
toolRes: {
citeLinks: [
{
name: 'Tool Ref',
url: 'https://example.com/ref'
},
{
name: 'Tool Ref',
url: 'https://example.com/ref'
}
]
},
errorText: 'tool failed'
}
},
{
teamId: testUser.teamId,
appId,
chatId,
chatItemDataId: aiDataId,
data: {
id: 'large-response',
moduleName: 'LLM',
moduleType: FlowNodeTypeEnum.chatNode,
historyPreview: 'large preview should not be projected',
toolRes: {
result: 'large result should not be projected'
}
}
}
]);
const result = await getChatItems({
appId,
chatId,
offset: 0,
limit: 10,
field: 'obj value responseData',
nodeResponseMode: 'preview'
});
const item = result.histories[0];
expect(item.responseData?.map((response) => response.id)).toEqual([
'dataset-response',
'tool-response',
'large-response'
]);
expect(item.responseData?.[0].quoteList?.map((quote) => quote.id)).toEqual([
citedQuoteId,
uncitedQuoteId
]);
expect(item.responseData?.[1].toolRes).toEqual({
citeLinks: [
{
name: 'Tool Ref',
url: 'https://example.com/ref'
},
{
name: 'Tool Ref',
url: 'https://example.com/ref'
}
]
});
expect(item.responseData?.[1].errorText).toBe('tool failed');
expect(item.responseData?.[2].historyPreview).toBeUndefined();
expect(item.responseData?.[2].toolRes?.result).toBeUndefined();
});
});
});
......
......@@ -4,10 +4,11 @@ import { ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
import { MongoChatItem } from '@fastgpt/service/core/chat/chatItemSchema';
import {
CHAT_DATA_ID_DUPLICATE_ERROR_MESSAGE,
assertNoExistingChatDataIds,
assertNoDuplicateChatDataIdsInRequest,
getChatMessagesDataIds,
validateChatRoundDataIds
} from '@fastgpt/service/core/chat/dataIdValidation';
} from '@fastgpt/service/core/chat/utils/dataIdValidation';
import { MongoApp } from '@fastgpt/service/core/app/schema';
import { AppTypeEnum } from '@fastgpt/global/core/app/constants';
import { getUser } from '@test/datas/users';
......@@ -58,7 +59,11 @@ describe('chat dataId validation', () => {
).toThrow(`${CHAT_DATA_ID_DUPLICATE_ERROR_MESSAGE}: history-1`);
});
it('should reject human and ai sharing one dataId', async () => {
it('should ignore empty values when checking duplicate dataIds in the current request', () => {
expect(() => assertNoDuplicateChatDataIdsInRequest([undefined, '', 'history-1'])).not.toThrow();
});
it('should allow human and ai sharing one round dataId', async () => {
await expect(
validateChatRoundDataIds({
appId,
......@@ -70,7 +75,32 @@ describe('chat dataId validation', () => {
},
responseChatItemId: 'same-data-id'
})
).rejects.toThrow(`${CHAT_DATA_ID_DUPLICATE_ERROR_MESSAGE}: same-data-id`);
).resolves.toBeUndefined();
});
it('should allow current round AI dataId when only an existing human item has the same dataId', async () => {
await MongoChatItem.create({
teamId: testUser.teamId,
tmbId: testUser.tmbId,
appId,
chatId,
dataId: 'same-round-id',
obj: ChatRoleEnum.Human,
value: [{ text: { content: 'old question' } }]
});
await expect(
validateChatRoundDataIds({
appId,
chatId,
userContent: {
obj: ChatRoleEnum.Human,
dataId: 'same-round-id',
value: [{ text: { content: 'hello' } }]
},
responseChatItemId: 'same-round-id'
})
).resolves.toBeUndefined();
});
it('should reject dataIds that already exist in chat items', async () => {
......@@ -113,6 +143,50 @@ describe('chat dataId validation', () => {
).resolves.toBeUndefined();
});
it('should skip current round AI dataId validation when responseChatItemId is missing', async () => {
await expect(
validateChatRoundDataIds({
appId,
chatId,
userContent: {
obj: ChatRoleEnum.Human,
dataId: 'human-only',
value: [{ text: { content: 'hello' } }]
}
})
).resolves.toBeUndefined();
});
it('should reject existing chat dataIds in generic history validation', async () => {
await MongoChatItem.create({
teamId: testUser.teamId,
tmbId: testUser.tmbId,
appId,
chatId,
dataId: 'existing-human',
obj: ChatRoleEnum.Human,
value: [{ text: { content: 'old question' } }]
});
await expect(
assertNoExistingChatDataIds({
appId,
chatId,
dataIds: ['missing', 'existing-human']
})
).rejects.toThrow(`${CHAT_DATA_ID_DUPLICATE_ERROR_MESSAGE}: existing-human`);
});
it('should skip generic history validation when all dataIds are empty', async () => {
await expect(
assertNoExistingChatDataIds({
appId,
chatId,
dataIds: [undefined, '']
})
).resolves.toBeUndefined();
});
it('should allow validating only the current human dataId for interactive submit', async () => {
await MongoChatItem.create({
teamId: testUser.teamId,
......
......@@ -75,7 +75,6 @@ describe('agent adapter useToolNodeResponse', () => {
nodeId: 'call_search',
moduleType: FlowNodeTypeEnum.tool,
moduleName: 'Search',
childTotalPoints: 999,
childrenResponses: [
{
id: 'existing_child',
......
......@@ -236,7 +236,6 @@ describe('runWorkflow node response persistence', () => {
})
];
const nodeResponseWriter = await createWorkflowNodeResponseWriter({
mode: 'replace',
teamId: '654a4107c32f3bf5f998452f',
appId,
chatId,
......@@ -427,7 +426,6 @@ describe('runWorkflow node response persistence', () => {
tmbId: '65ab7007462ada7dbb899948'
};
const nodeResponseWriter = await createWorkflowNodeResponseWriter({
mode: 'replace',
teamId: runningAppInfo.teamId,
appId,
chatId,
......@@ -505,12 +503,224 @@ describe('runWorkflow node response persistence', () => {
}
});
it('returns the runtime node response id when workflow pauses on an interactive node', async () => {
const appId = '67e0d5535c02d1d5cdede723';
const chatId = 'workflow-interactive-node-response-chat';
const responseChatItemId = 'workflow-interactive-node-response-ai-item';
const runningAppInfo = {
id: appId,
name: 'Workflow Interactive App',
teamId: '654a4107c32f3bf5f998452f',
tmbId: '65ab7007462ada7dbb899948'
};
const runtimeNodes: RuntimeNodeItemType[] = [
makeNode('select_node', FlowNodeTypeEnum.userSelect, {
isEntry: true,
inputs: [
makeInput({
key: NodeInputKeyEnum.description,
value: 'Choose one',
valueType: WorkflowIOValueTypeEnum.string
}),
makeInput({
key: NodeInputKeyEnum.userSelectOptions,
value: [
{
key: 'a',
value: 'A'
}
],
valueType: WorkflowIOValueTypeEnum.arrayObject
})
]
})
];
const nodeResponseWriter = await createWorkflowNodeResponseWriter({
teamId: runningAppInfo.teamId,
appId,
chatId,
chatItemDataId: responseChatItemId
});
const result = await runWorkflow({
apiVersion: 'v2',
mode: 'chat',
chatId,
responseChatItemId,
runningAppInfo,
runningUserInfo: {
teamId: runningAppInfo.teamId,
tmbId: runningAppInfo.tmbId,
teamName: 'team',
memberName: 'member',
contact: '',
username: 'user'
},
uid: 'test-user',
lang: 'zh-CN',
histories: [],
query: [{ type: 'text', text: { content: 'pause' } }],
variables: {},
chatConfig: {},
runtimeNodes,
runtimeEdges: [],
variableState: await WorkflowVariableState.create({
timezone: 'Asia/Shanghai',
runningAppInfo,
uid: 'test-user',
chatId,
responseChatItemId,
histories: [],
variablesConfig: [],
inputVariables: {},
externalVariables: {}
}),
externalProvider: {},
workflowDispatchDeep: 0,
maxRunTimes: 20,
stream: false,
responseAllData: true,
responseDetail: true,
nodeResponseWriter,
checkIsStopping: () => false
} as any);
await nodeResponseWriter.close();
const interactive = result.workflowInteractiveResponse;
expect(interactive).toEqual(
expect.objectContaining({
type: 'userSelect',
entryNodeIds: ['select_node'],
nodeResponseId: expect.any(String)
})
);
const rows = await MongoChatItemResponse.find({
appId,
chatId,
chatItemDataId: responseChatItemId
}).lean();
expect(rows).toHaveLength(0);
});
it('keeps nested interactive nodeResponseIds scoped to their own layer', async () => {
const originalToolCallDispatch = callbackMap[FlowNodeTypeEnum.toolCall];
const childNodeResponseId = 'child_select_response';
callbackMap[FlowNodeTypeEnum.toolCall] = vi.fn(async () => ({
[DispatchNodeResponseKeyEnum.interactive]: {
type: 'toolChildrenInteractive',
params: {
childrenResponse: {
type: 'userSelect',
nodeResponseId: childNodeResponseId,
entryNodeIds: ['select_node'],
memoryEdges: [],
nodeOutputs: [],
params: {
description: 'Choose one',
userSelectOptions: []
}
},
toolParams: {
memoryRequestMessages: [],
toolCallId: 'call_search'
}
}
},
[DispatchNodeResponseKeyEnum.nodeResponse]: {
totalPoints: 0
}
}));
try {
const appId = '67e0d5535c02d1d5cdede724';
const chatId = 'workflow-nested-interactive-node-response-chat';
const responseChatItemId = 'workflow-nested-interactive-node-response-ai-item';
const runningAppInfo = {
id: appId,
name: 'Workflow Nested Interactive App',
teamId: '654a4107c32f3bf5f998452f',
tmbId: '65ab7007462ada7dbb899948'
};
const runtimeNodes: RuntimeNodeItemType[] = [
makeNode('tool_call_node', FlowNodeTypeEnum.toolCall, {
isEntry: true
})
];
const nodeResponseWriter = await createWorkflowNodeResponseWriter({
teamId: runningAppInfo.teamId,
appId,
chatId,
chatItemDataId: responseChatItemId
});
const result = await runWorkflow({
apiVersion: 'v2',
mode: 'chat',
chatId,
responseChatItemId,
runningAppInfo,
runningUserInfo: {
teamId: runningAppInfo.teamId,
tmbId: runningAppInfo.tmbId,
teamName: 'team',
memberName: 'member',
contact: '',
username: 'user'
},
uid: 'test-user',
lang: 'zh-CN',
histories: [],
query: [{ type: 'text', text: { content: 'pause' } }],
variables: {},
chatConfig: {},
runtimeNodes,
runtimeEdges: [],
variableState: await WorkflowVariableState.create({
timezone: 'Asia/Shanghai',
runningAppInfo,
uid: 'test-user',
chatId,
responseChatItemId,
histories: [],
variablesConfig: [],
inputVariables: {},
externalVariables: {}
}),
externalProvider: {},
workflowDispatchDeep: 0,
maxRunTimes: 20,
stream: false,
responseAllData: true,
responseDetail: true,
nodeResponseWriter,
checkIsStopping: () => false
} as any);
await nodeResponseWriter.close();
const interactive = result.workflowInteractiveResponse;
expect(interactive).toEqual(
expect.objectContaining({
type: 'toolChildrenInteractive',
entryNodeIds: ['tool_call_node'],
nodeResponseId: expect.any(String),
params: expect.objectContaining({
childrenResponse: expect.objectContaining({
nodeResponseId: childNodeResponseId
})
})
})
);
expect(interactive?.nodeResponseId).not.toBe(childNodeResponseId);
} finally {
callbackMap[FlowNodeTypeEnum.toolCall] = originalToolCallDispatch;
}
});
it('streams batched node responses into v2 flat rows and composes childrenResponses on detail read', async () => {
const loopItems = ['alpha', 'beta', 'gamma', 'delta', 'epsilon', 'zeta'];
const { runtimeNodes, runtimeEdges } = createLoopRunWorkflow(loopItems);
const responseEvents: unknown[] = [];
const nodeResponseWriter = await createWorkflowNodeResponseWriter({
mode: 'replace',
teamId: '654a4107c32f3bf5f998452f',
appId: '67e0d5535c02d1d5cdede71f',
chatId: 'workflow-persistence-chat',
......
......@@ -3,10 +3,7 @@ import { EventEmitter } from 'node:events';
import { createServer } from 'node:http';
import type { AddressInfo } from 'node:net';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import {
WorkflowQueue,
mergeAssistantResponseAnswerText
} from '@fastgpt/service/core/workflow/dispatch/index';
import { WorkflowQueue } from '@fastgpt/service/core/workflow/dispatch/index';
import { getWorkflowNodeRunParams } from '@fastgpt/service/core/workflow/dispatch/utils/runtime';
import { createClientAbortTracker } from '@fastgpt/service/core/workflow/dispatch/utils/clientAbort';
import { createNode, createEdge } from '../utils';
......@@ -40,41 +37,6 @@ const waitWithTimeout = async <T>(promise: Promise<T>, timeoutMs: number, label:
}
};
describe('mergeAssistantResponseAnswerText', () => {
it('should merge consecutive plain text assistant values', () => {
expect(
mergeAssistantResponseAnswerText([
{ text: { content: 'first' } },
{ text: { content: ' second' } }
])
).toEqual([{ text: { content: 'first second' } }]);
});
it('should preserve reasoning boundaries for consecutive AI chat nodes', () => {
expect(
mergeAssistantResponseAnswerText([
{
reasoning: { content: 'think 1' },
text: { content: 'answer 1' }
},
{
reasoning: { content: 'think 2' },
text: { content: 'answer 2' }
}
])
).toEqual([
{
reasoning: { content: 'think 1' },
text: { content: 'answer 1' }
},
{
reasoning: { content: 'think 2' },
text: { content: 'answer 2' }
}
]);
});
});
describe('createClientAbortTracker', () => {
const mockRes = (overrides: Record<string, any> = {}) => {
const res = new EventEmitter() as any;
......
......@@ -425,7 +425,6 @@ describe('runLoopRun (integration with mocked runWorkflow)', () => {
);
const nodeResponse = result[DispatchNodeResponseKeyEnum.nodeResponse];
expect(nodeResponse.errorText).toBe('workflow:loop_run_conditional_requires_break');
expect(nodeResponse.mergeSignId).toBe('loopRun1');
expect(runWorkflowMock).not.toHaveBeenCalled();
});
......@@ -542,7 +541,7 @@ describe('runLoopRun (integration with mocked runWorkflow)', () => {
expect(history[1]).toMatchObject({ iteration: 2, success: true });
});
it('lastInteractive 恢复 → 只累计恢复后的增量 child 统计', async () => {
it('lastInteractive 恢复 → 完成时只写入恢复后的 wrapper 增量统计', async () => {
const nodeResponseWriter = {
recordWithParent: vi.fn().mockResolvedValue([])
};
......@@ -602,8 +601,8 @@ describe('runLoopRun (integration with mocked runWorkflow)', () => {
expect(nodeResponse.childResponseCount).toBe(3);
expect(nodeResponseWriter.recordWithParent).toHaveBeenCalledTimes(1);
expect(nodeResponseWriter.recordWithParent.mock.calls[0][0][0]).toMatchObject({
id: 'loopRun1_iter_1',
totalPoints: 6,
id: 'loop-parent-response:iter:1',
totalPoints: 5,
childResponseCount: 2
});
});
......@@ -791,11 +790,13 @@ describe('runLoopRun (integration with mocked runWorkflow)', () => {
const result: any = await dispatchLoopRun(props);
const nodeResponse = result[DispatchNodeResponseKeyEnum.nodeResponse];
expect(runWorkflowMock.mock.calls[0][0].nodeResponseParentId).toBe('loopRun1_iter_1');
expect(runWorkflowMock.mock.calls[0][0].nodeResponseParentId).toBe(
'loop-parent-response:iter:1'
);
expect(nodeResponseWriter.recordWithParent).toHaveBeenCalledTimes(1);
expect(nodeResponseWriter.recordWithParent.mock.calls[0][1]).toBe('loop-parent-response');
expect(nodeResponseWriter.recordWithParent.mock.calls[0][0][0]).toMatchObject({
id: 'loopRun1_iter_1',
id: 'loop-parent-response:iter:1',
childResponseCount: 2
});
expect(
......@@ -849,7 +850,7 @@ describe('runLoopRun (integration with mocked runWorkflow)', () => {
});
});
it('interactive 中断轮:不内嵌 partial loopRunDetail,父响应保留 child 统计', async () => {
it('interactive 中断轮:写入本轮 wrapper,确保已写 child 可挂到 loop 下', async () => {
const interactivePayload: any = {
entryNodeIds: ['userSelectNode'],
memoryEdges: [],
......@@ -864,17 +865,30 @@ describe('runLoopRun (integration with mocked runWorkflow)', () => {
})
)
);
const nodeResponseWriter = {
recordWithParent: vi.fn().mockResolvedValue([])
};
const props = makeProps({
[NodeInputKeyEnum.loopRunMode]: LoopRunModeEnum.array,
[NodeInputKeyEnum.loopRunInputArray]: ['a', 'b'],
[NodeInputKeyEnum.childrenNodeIdList]: ['startNode', 'chatNode']
});
const props = {
...makeProps({
[NodeInputKeyEnum.loopRunMode]: LoopRunModeEnum.array,
[NodeInputKeyEnum.loopRunInputArray]: ['a', 'b'],
[NodeInputKeyEnum.childrenNodeIdList]: ['startNode', 'chatNode']
}),
nodeResponseWriter,
nodeResponseParentId: 'loop-parent-response'
};
const result: any = await dispatchLoopRun(props);
const nodeResponse = result[DispatchNodeResponseKeyEnum.nodeResponse];
expect(nodeResponse.loopRunDetail).toBeUndefined();
expect(nodeResponse.childResponseCount).toBe(2);
expect(nodeResponseWriter.recordWithParent).toHaveBeenCalledTimes(1);
expect(nodeResponseWriter.recordWithParent.mock.calls[0][1]).toBe('loop-parent-response');
expect(nodeResponseWriter.recordWithParent.mock.calls[0][0][0]).toMatchObject({
id: 'loop-parent-response:iter:1',
childResponseCount: 1
});
});
it('array mode 数组长度 === max → 跑满且不报超限(回归:== max 不算超限)', async () => {
......
......@@ -51,7 +51,12 @@ vi.mock('@fastgpt/service/core/workflow/dispatch', () => ({
}));
vi.mock('@fastgpt/service/core/chat/saveChat', () => ({
pushChatRecords: vi.fn()
finalizeChatRound: vi.fn(),
failChatRound: vi.fn()
}));
vi.mock('@fastgpt/service/core/chat/utils/prepare', () => ({
preChatRound: vi.fn()
}));
beforeEach(() => {
......
import { beforeEach, describe, expect, it, vi } from 'vitest';
import { ChatRoleEnum } from '@fastgpt/global/core/chat/constants';
import { PublishChannelEnum } from '@fastgpt/global/support/outLink/constant';
import { outlinkInvokeChat } from '@fastgpt/service/support/outLink/runtime/utils';
import { dispatchWorkFlow } from '@fastgpt/service/core/workflow/dispatch';
import { getAppLatestVersion } from '@fastgpt/service/core/app/version/controller';
import { MongoApp } from '@fastgpt/service/core/app/schema';
import { getChatItems } from '@fastgpt/service/core/chat/controller';
import { preChatRound } from '@fastgpt/service/core/chat/utils/prepare';
import { failChatRound, finalizeChatRound } from '@fastgpt/service/core/chat/saveChat';
import { authOutLinkLimit } from '@fastgpt/service/support/outLink/runtime/auth';
import { addOutLinkUsage } from '@fastgpt/service/support/outLink/tools';
import { getRunningUserInfoByTmbId } from '@fastgpt/service/support/user/team/utils';
vi.mock('@fastgpt/service/core/app/schema', () => ({
MongoApp: {
findById: vi.fn()
}
}));
vi.mock('@fastgpt/service/core/app/version/controller', () => ({
getAppLatestVersion: vi.fn()
}));
vi.mock('@fastgpt/service/core/chat/controller', () => ({
getChatItems: vi.fn()
}));
vi.mock('@fastgpt/service/core/chat/chatSchema', () => ({
MongoChat: {
findOne: vi.fn(() => ({
lean: () => ({ variables: { retained: 'value' } })
})),
updateOne: vi.fn()
}
}));
vi.mock('@fastgpt/service/core/chat/chatItemSchema', () => ({
MongoChatItem: {
updateMany: vi.fn()
}
}));
vi.mock('@fastgpt/service/core/chat/utils/prepare', () => ({
preChatRound: vi.fn()
}));
vi.mock('@fastgpt/service/core/chat/saveChat', () => ({
finalizeChatRound: vi.fn(),
failChatRound: vi.fn()
}));
vi.mock('@fastgpt/service/core/chat/chatGenerateStatus', () => ({
updateChatGenerateStatus: vi.fn()
}));
vi.mock('@fastgpt/service/core/workflow/dispatch', () => ({
dispatchWorkFlow: vi.fn()
}));
vi.mock('@fastgpt/service/support/user/team/utils', () => ({
getRunningUserInfoByTmbId: vi.fn()
}));
vi.mock('@fastgpt/service/support/outLink/runtime/auth', () => ({
authOutLinkLimit: vi.fn()
}));
vi.mock('@fastgpt/service/support/outLink/tools', () => ({
addOutLinkUsage: vi.fn()
}));
vi.mock('@fastgpt/global/core/workflow/runtime/utils', async (importOriginal) => {
const actual =
await importOriginal<typeof import('@fastgpt/global/core/workflow/runtime/utils')>();
return {
...actual,
getMaxHistoryLimitFromNodes: vi.fn(() => 10),
getWorkflowEntryNodeIds: vi.fn(() => ['start']),
storeNodes2RuntimeNodes: vi.fn(() => [{ nodeId: 'runtime-start' }]),
storeEdges2RuntimeEdges: vi.fn(() => [])
};
});
describe('outlinkInvokeChat', () => {
const outLinkConfig = {
_id: 'outlink-id',
shareId: 'share-id',
teamId: 'team-id',
tmbId: 'tmb-id',
appId: 'app-id',
name: 'OutLink',
usagePoints: 0,
lastTime: new Date('2026-06-14T00:00:00.000Z'),
type: PublishChannelEnum.feishu,
showCite: true,
showRunningStatus: true,
showSkillReferences: false,
showFullText: true,
canDownloadSource: true,
showWholeResponse: true,
app: undefined
};
const query = [{ text: { content: 'hello outlink' } }];
beforeEach(() => {
vi.clearAllMocks();
vi.mocked(MongoApp.findById).mockReturnValue({
lean: () => ({
_id: 'app-id',
name: 'App',
teamId: 'app-team-id',
tmbId: 'app-tmb-id'
})
} as any);
vi.mocked(getAppLatestVersion).mockResolvedValue({
nodes: [{ nodeId: 'start', inputs: [], outputs: [] }],
edges: [],
chatConfig: { variables: [] }
} as any);
vi.mocked(getChatItems).mockResolvedValue({
histories: []
} as any);
vi.mocked(authOutLinkLimit).mockResolvedValue(undefined as any);
vi.mocked(getRunningUserInfoByTmbId).mockResolvedValue({
teamId: 'team-id',
tmbId: 'tmb-id'
} as any);
vi.mocked(preChatRound).mockResolvedValue({
chatId: 'prepared-chat-id',
responseChatItemId: 'message-id',
shouldPersistChatRound: true,
shouldFinalizePreparedRound: true
});
vi.mocked(dispatchWorkFlow).mockResolvedValue({
assistantResponses: [{ text: { content: 'answer' } }],
newVariables: { next: 'value' },
flowUsages: [{ totalPoints: 3 }],
durationSeconds: 1.5,
system_memories: { memory: 'value' },
nodeResponseSummary: {
citeCollectionIds: [],
errorCount: 0,
totalPoints: 3
}
} as any);
vi.mocked(finalizeChatRound).mockResolvedValue(undefined as any);
vi.mocked(failChatRound).mockResolvedValue(undefined as any);
vi.mocked(addOutLinkUsage).mockResolvedValue(undefined as any);
});
it('prepares and finalizes outlink chat rounds using messageId as the round dataId', async () => {
const onReply = vi.fn().mockResolvedValue(undefined);
await outlinkInvokeChat({
outLinkConfig,
chatId: 'chat-id',
query,
messageId: 'message-id',
chatUserId: 'chat-user-id',
onReply
});
expect(preChatRound).toHaveBeenCalledWith(
expect.objectContaining({
appId: 'app-id',
chatId: 'chat-id',
sourceName: 'OutLink',
shareId: 'share-id',
outLinkUid: 'chat-user-id',
responseChatItemId: 'message-id',
userContent: {
dataId: 'message-id',
obj: ChatRoleEnum.Human,
value: query
}
})
);
expect(dispatchWorkFlow).toHaveBeenCalledWith(
expect.objectContaining({
chatId: 'prepared-chat-id',
responseChatItemId: 'message-id'
})
);
expect(finalizeChatRound).toHaveBeenCalledWith(
expect.objectContaining({
chatId: 'prepared-chat-id',
aiContent: expect.objectContaining({
dataId: 'message-id',
value: [{ text: { content: 'answer' } }]
})
})
);
expect(onReply).toHaveBeenCalledWith('answer');
expect(addOutLinkUsage).toHaveBeenCalledWith({
shareId: 'share-id',
totalPoints: 3
});
});
it('marks prepared outlink round as failed when workflow dispatch rejects', async () => {
const error = new Error('workflow failed');
vi.mocked(dispatchWorkFlow).mockRejectedValue(error);
await outlinkInvokeChat({
outLinkConfig,
chatId: 'chat-id',
query,
messageId: 'message-id',
chatUserId: 'chat-user-id',
onReply: vi.fn().mockResolvedValue(undefined)
});
expect(finalizeChatRound).not.toHaveBeenCalled();
expect(failChatRound).toHaveBeenCalledWith({
appId: 'app-id',
chatId: 'prepared-chat-id',
responseChatItemId: 'message-id',
error
});
});
});
Subproject commit 298aa5c17c853e98ff15672721acb343b3e8b365
Subproject commit 8c76261acbe8a406cedd8c55e306cbac66b90b7e
{
"name": "@fastgpt/app",
"version": "4.15.0-4",
"version": "4.15.0-5",
"private": false,
"browserslist": [
"Chrome >= 80",
......
.waitingAnimation> :last-child::after {
display: inline-block;
content: '';
width: 3px;
max-width: 3px;
height: 14px;
transform: translate(4px, 2px) scaleY(1.3);
background-color: var(--chakra-colors-primary-700);
animation: blink 0.6s infinite;
vertical-align: baseline;
}
.waitingAnimation > .htmlCodeBlock:last-child::after {
display: none;
.waitingAnimation {
:global(.stream-char) {
opacity: 0;
filter: blur(1px);
transform: translateY(1px);
animation: streamFadeIn 420ms cubic-bezier(0.16, 1, 0.3, 1) forwards;
}
:global(.katex-display) :global(.katex-html) span {
animation: none !important;
}
}
.animation {
......@@ -28,7 +26,6 @@
}
@keyframes blink {
from,
to {
opacity: 0;
......@@ -39,11 +36,25 @@
}
}
.markdown>*:first-child {
@keyframes streamFadeIn {
from {
opacity: 0;
filter: blur(1px);
transform: translateY(1px);
}
to {
opacity: 1;
filter: blur(0);
transform: translateY(0);
}
}
.markdown > *:first-child {
margin-top: 0 !important;
}
.markdown>*:last-child {
.markdown > *:last-child {
margin-bottom: 0 !important;
}
......@@ -158,13 +169,13 @@
margin: 14px 0;
}
.markdown>h2:first-child,
.markdown>h1:first-child,
.markdown>h1:first-child+h2,
.markdown>h3:first-child,
.markdown>h4:first-child,
.markdown>h5:first-child,
.markdown>h6:first-child {
.markdown > h2:first-child,
.markdown > h1:first-child,
.markdown > h1:first-child + h2,
.markdown > h3:first-child,
.markdown > h4:first-child,
.markdown > h5:first-child,
.markdown > h6:first-child {
margin-top: 0;
padding-top: 0;
}
......@@ -179,12 +190,12 @@
padding-top: 0;
}
.markdown h1+p,
.markdown h2+p,
.markdown h3+p,
.markdown h4+p,
.markdown h5+p,
.markdown h6+p {
.markdown h1 + p,
.markdown h2 + p,
.markdown h3 + p,
.markdown h4 + p,
.markdown h5 + p,
.markdown h6 + p {
margin-top: 0;
}
......@@ -208,8 +219,8 @@
padding-left: 0.25em;
}
.markdown ul li>*:first-child,
.markdown ol li>*:first-child {
.markdown ul li > *:first-child,
.markdown ol li > *:first-child {
margin-top: 0;
}
......@@ -229,11 +240,11 @@
padding: 0;
}
.markdown dl dt>*:first-child {
.markdown dl dt > *:first-child {
margin-top: 0;
}
.markdown dl dt>*:last-child {
.markdown dl dt > *:last-child {
margin-bottom: 0;
}
......@@ -242,11 +253,11 @@
padding: 0 15px;
}
.markdown dl dd>*:first-child {
.markdown dl dd > *:first-child {
margin-top: 0;
}
.markdown dl dd>*:last-child {
.markdown dl dd > *:last-child {
margin-bottom: 0;
}
......@@ -256,11 +267,11 @@
padding: 0 15px;
}
.markdown blockquote>*:first-child {
.markdown blockquote > *:first-child {
margin-top: 0;
}
.markdown blockquote>*:last-child {
.markdown blockquote > *:last-child {
margin-bottom: 0;
}
......@@ -295,7 +306,7 @@
overflow: hidden;
}
.markdown span.frame>span {
.markdown span.frame > span {
border: 1px solid #dddddd;
display: block;
float: left;
......@@ -323,7 +334,7 @@
overflow: hidden;
}
.markdown span.align-center>span {
.markdown span.align-center > span {
display: block;
margin: 13px auto 0;
overflow: hidden;
......@@ -341,7 +352,7 @@
overflow: hidden;
}
.markdown span.align-right>span {
.markdown span.align-right > span {
display: block;
margin: 13px 0 0;
overflow: hidden;
......@@ -371,7 +382,7 @@
overflow: hidden;
}
.markdown span.float-right>span {
.markdown span.float-right > span {
display: block;
margin: 13px auto 0;
overflow: hidden;
......@@ -387,7 +398,7 @@
padding: 0 5px;
}
.markdown pre>code {
.markdown pre > code {
background: none repeat scroll 0 0 transparent;
border: medium none;
margin: 0;
......
......@@ -11,10 +11,11 @@ import styles from './index.module.scss';
import dynamic from 'next/dynamic';
import { Box } from '@chakra-ui/react';
import { CodeClassNameEnum, mdTextFormat } from './utils';
import { CodeClassNameEnum, hideStreamingIncompleteMarkdownTail, mdTextFormat } from './utils';
import type { AProps } from './A';
import MarkdownTable from '@fastgpt/web/components/common/Markdown/MarkdownTable';
import { MarkdownRendererRuntimeContext } from './runtimeContext';
import { useStreamAnimatedRehypePlugin } from './rehypeStreamAnimated';
const CodeLight = dynamic(() => import('./codeBlock/CodeLight'), { ssr: false });
const MermaidCodeBlock = dynamic(() => import('./img/MermaidCodeBlock'), { ssr: false });
......@@ -80,6 +81,7 @@ type Props = {
className?: string;
autoPreviewHtmlCodeBlock?: boolean;
} & AProps;
const Markdown = (props: Props) => {
const source = props.source || '';
......@@ -112,10 +114,20 @@ const MarkdownRender = ({
);
const formatSource = useMemo(() => {
if (showAnimation || forbidZhFormat) return source;
if (showAnimation) return hideStreamingIncompleteMarkdownTail(source);
if (forbidZhFormat) return source;
return mdTextFormat(source);
}, [forbidZhFormat, showAnimation, source]);
const streamAnimatedRehypePlugin = useStreamAnimatedRehypePlugin();
const rehypePlugins = useMemo(
() =>
showAnimation
? [RehypeKatex, [RehypeExternalLinks, { target: '_blank' }], streamAnimatedRehypePlugin]
: [RehypeKatex, [RehypeExternalLinks, { target: '_blank' }]],
[showAnimation, streamAnimatedRehypePlugin]
);
const urlTransform = useCallback((val: string) => {
return val;
}, []);
......@@ -129,7 +141,7 @@ const MarkdownRender = ({
${showAnimation ? `${formatSource ? styles.waitingAnimation : styles.animation}` : ''}
`}
remarkPlugins={[RemarkMath, [RemarkGfm, { singleTilde: false }], RemarkBreaks]}
rehypePlugins={[RehypeKatex, [RehypeExternalLinks, { target: '_blank' }]]}
rehypePlugins={rehypePlugins as any}
components={markdownComponents}
urlTransform={urlTransform}
>
......
import { useMemo } from 'react';
const STREAM_ANIMATED_BLOCK_TAGS = new Set(['p', 'h1', 'h2', 'h3', 'h4', 'h5', 'h6']);
const STREAM_ANIMATED_SKIP_TAGS = new Set(['pre', 'code', 'table', 'svg']);
const STREAM_ANIMATED_LIST_INLINE_TAGS = new Set([
'a',
'abbr',
'b',
'cite',
'del',
'em',
'i',
'ins',
'kbd',
'mark',
's',
'small',
'span',
'strong',
'sub',
'sup',
'u'
]);
type HastElement = {
type: 'element';
tagName: string;
properties?: Record<string, any>;
children?: HastNode[];
};
type HastNode =
| HastElement
| {
type: 'text';
value: string;
}
| {
type: string;
[key: string]: any;
};
type HastRoot = {
type: 'root';
children?: HastNode[];
};
/**
* 流式输出时按 Lobe UI 的方式把新增文本拆成字符 span 做淡入动画。
* 只处理段落、标题和列表里的文本,跳过代码块、表格、公式等高复杂度内容,避免流式阶段 DOM 过重。
*/
export const rehypeStreamAnimated = ({ className = 'stream-char' }: { className?: string }) => {
return (tree: HastRoot) => {
const isHastElement = (node: HastNode): node is HastElement => {
return node.type === 'element' && typeof (node as HastElement).tagName === 'string';
};
const hasClass = (node: HastElement, cls: string) => {
const className = node.properties?.className;
if (Array.isArray(className)) return className.some((item) => String(item).includes(cls));
if (typeof className === 'string') return className.includes(cls);
return false;
};
const shouldSkip = (node: HastElement) => {
return STREAM_ANIMATED_SKIP_TAGS.has(node.tagName) || hasClass(node, 'katex');
};
const createStreamCharNode = (char: string): HastElement => {
return {
type: 'element',
tagName: 'span',
properties: { className },
children: [{ type: 'text', value: char }]
};
};
const wrapText = (node: HastElement) => {
const newChildren: HastNode[] = [];
node.children?.forEach((child) => {
if (child.type === 'text') {
for (const char of child.value) {
newChildren.push(createStreamCharNode(char));
}
return;
}
if (isHastElement(child) && !shouldSkip(child)) {
wrapText(child);
}
newChildren.push(child);
});
node.children = newChildren;
};
const wrapListItemText = (node: HastElement) => {
const newChildren: HastNode[] = [];
node.children?.forEach((child) => {
if (child.type === 'text') {
for (const char of child.value) {
newChildren.push(createStreamCharNode(char));
}
return;
}
if (
isHastElement(child) &&
(child.tagName === 'p' || STREAM_ANIMATED_LIST_INLINE_TAGS.has(child.tagName))
) {
wrapText(child);
}
newChildren.push(child);
});
node.children = newChildren;
};
const visitElement = (node: HastNode): 'skip' | undefined => {
if (!isHastElement(node)) return;
if (shouldSkip(node)) return 'skip';
if (node.tagName === 'li') {
wrapListItemText(node);
return 'skip';
}
if (STREAM_ANIMATED_BLOCK_TAGS.has(node.tagName)) {
wrapText(node);
return 'skip';
}
node.children?.forEach((child) => {
if (visitElement(child) === 'skip') return;
});
};
tree.children?.forEach((child) => {
visitElement(child);
});
};
};
/**
* 返回流式 Markdown 字符淡入插件。
* 不在 React render 期间读写 ref;依赖 React 复用已存在的字符 span,仅让新增字符节点触发 CSS animation。
*/
export const useStreamAnimatedRehypePlugin = () => {
return useMemo(() => [rehypeStreamAnimated, { className: 'stream-char' }], []);
};
......@@ -14,6 +14,188 @@ export enum CodeClassNameEnum {
audio = 'audio'
}
const streamingIncompleteMarkdownTailPatterns = [
/!\[[^\]\n]*\]\([^\s\n)]*$/,
/!\[[^\]\n]*\]$/,
/!\[[^\]\n]*$/,
/\[[^\]\n]*\]\([^\s\n)]*$/
];
const streamingIncompleteTextMarkdownTailMarkers = ['**', '__', '~~'] as const;
const streamingIncompleteItalicMarkdownTailMarkers = ['*', '_'] as const;
const STREAMING_INCOMPLETE_MEDIA_MARKDOWN_TAIL_MAX_HIDE_LENGTH = 100;
const STREAMING_INCOMPLETE_TEXT_MARKDOWN_TAIL_MAX_HIDE_LENGTH = 50;
/**
* 流式输出时隐藏尾部尚未闭合的 Markdown 片段。
*
* 该函数只处理“还在尾部继续生长”的候选语法,不补全括号,也不处理正文中间的
* 不完整 Markdown。这样图片不会因为半截 URL 被提前请求;如果后续输出出现空白、
* 换行等打断符,或候选片段过长,候选语法会按原文重新显示。
*/
export const hideStreamingIncompleteMarkdownTail = (text: string) => {
const isEscapedAt = (text: string, index: number) => {
let slashCount = 0;
for (let i = index - 1; i >= 0 && text[i] === '\\'; i -= 1) {
slashCount += 1;
}
return slashCount % 2 === 1;
};
const isInsideInlineCode = (text: string, index: number) => {
let backtickCount = 0;
for (let i = 0; i < index; i += 1) {
if (text[i] !== '`' || isEscapedAt(text, i)) continue;
if (text.substring(i, i + 3) === '```') {
i += 2;
continue;
}
backtickCount += 1;
}
return backtickCount % 2 === 1;
};
const isInsideOpenCodeFence = (text: string) => {
const fences = text.match(/(^|\n)```/g);
return !!fences && fences.length % 2 === 1;
};
const isWordChar = (char?: string) => {
return !!char && /[\p{L}\p{N}_]/u.test(char);
};
const isWhitespace = (char?: string) => {
return !!char && /\s/.test(char);
};
const isCodeFenceBacktickAt = (text: string, index: number) => {
return (
text.substring(index, index + 3) === '```' ||
text.substring(index - 1, index + 2) === '```' ||
text.substring(index - 2, index + 1) === '```'
);
};
const getUnescapedMarkerPositions = (text: string, marker: string) => {
const positions: number[] = [];
for (let i = 0; i <= text.length - marker.length; i += 1) {
if (text.substring(i, i + marker.length) !== marker) continue;
if (isEscapedAt(text, i)) continue;
if (marker === '`' && isCodeFenceBacktickAt(text, i)) continue;
if (marker !== '`' && isInsideInlineCode(text, i)) continue;
positions.push(i);
i += marker.length - 1;
}
return positions;
};
const isSingleItalicMarkerAt = (text: string, index: number, marker: '*' | '_') => {
if (text[index] !== marker) return false;
if (isEscapedAt(text, index)) return false;
if (isInsideInlineCode(text, index)) return false;
if (text[index - 1] === marker || text[index + 1] === marker) return false;
return true;
};
const getUnescapedSingleItalicMarkerPositions = (text: string, marker: '*' | '_') => {
const positions: number[] = [];
for (let i = 0; i < text.length; i += 1) {
if (!isSingleItalicMarkerAt(text, i, marker)) continue;
positions.push(i);
}
return positions;
};
const shouldHideStreamingTextMarkdownTail = ({
text,
startIndex,
marker
}: {
text: string;
startIndex: number;
marker: string;
}) => {
if (text.length - startIndex > STREAMING_INCOMPLETE_TEXT_MARKDOWN_TAIL_MAX_HIDE_LENGTH) {
return false;
}
const nextChar = text[startIndex + marker.length];
if (isWhitespace(nextChar)) return false;
// 双标记容易和普通文本粘连,要求左侧不是单词字符,降低误隐藏概率。
if (marker !== '`' && isWordChar(text[startIndex - 1])) return false;
return true;
};
const getStreamingIncompleteTextMarkdownTailStart = (text: string) => {
const inlineCodePositions = getUnescapedMarkerPositions(text, '`');
if (inlineCodePositions.length % 2 === 1) {
const startIndex = inlineCodePositions[inlineCodePositions.length - 1];
if (shouldHideStreamingTextMarkdownTail({ text, startIndex, marker: '`' })) {
return startIndex;
}
}
for (const marker of streamingIncompleteTextMarkdownTailMarkers) {
const positions = getUnescapedMarkerPositions(text, marker);
if (positions.length % 2 === 0) continue;
const startIndex = positions[positions.length - 1];
if (shouldHideStreamingTextMarkdownTail({ text, startIndex, marker })) {
return startIndex;
}
}
for (const marker of streamingIncompleteItalicMarkdownTailMarkers) {
const positions = getUnescapedSingleItalicMarkerPositions(text, marker);
if (positions.length % 2 === 0) continue;
const startIndex = positions[positions.length - 1];
if (shouldHideStreamingTextMarkdownTail({ text, startIndex, marker })) {
return startIndex;
}
}
};
if (!text || isInsideOpenCodeFence(text)) return text;
for (const pattern of streamingIncompleteMarkdownTailPatterns) {
const match = text.match(pattern);
const startIndex = match?.index;
if (startIndex === undefined) continue;
if (text[startIndex] === '[' && text[startIndex - 1] === '!') continue;
if (isEscapedAt(text, startIndex)) continue;
if (isInsideInlineCode(text, startIndex)) continue;
if (text.length - startIndex > STREAMING_INCOMPLETE_MEDIA_MARKDOWN_TAIL_MAX_HIDE_LENGTH) {
continue;
}
return text.slice(0, startIndex);
}
const textMarkdownStartIndex = getStreamingIncompleteTextMarkdownTailStart(text);
if (textMarkdownStartIndex !== undefined) {
return text.slice(0, textMarkdownStartIndex);
}
return text;
};
export const mdTextFormat = (text: string) => {
// 处理 Windows 文件路径中的反斜杠,防止被 Markdown 转义:C:\path\file 或 c:\path\file
text = text.replace(/([A-Za-z]:\\[^\s`\[\]()]*)/g, (match) => {
......
......@@ -19,6 +19,7 @@ import AIChatBubble, { shouldFilterAiValue } from './AIChatBubble';
import type { ChatBoxInputType } from '../type';
import { hasAiAnswerContent } from './AIChatBubble/utils';
import ChatErrorCard from './ChatErrorCard';
import { shouldShowChatItemInlineError } from '../utils/error';
const colorMap = {
[ChatStatusEnum.loading]: {
......@@ -92,6 +93,11 @@ const ChatItem = (props: Props) => {
message: t(errorText?.errorText || chat.errorMsg || 'Unknow error')
};
}, [chat.errorMsg, chat.moduleName, errorText, t]);
const showInlineError = shouldShowChatItemInlineError({
hasInlineError: !!inlineErrorInfo,
isChatting,
isLastChild
});
const isChatLog = chatType === 'log';
......@@ -286,7 +292,7 @@ const ChatItem = (props: Props) => {
i === splitAiResponseResults.length - 1 ? (
<>
{/* error message */}
{inlineErrorInfo && (
{showInlineError && inlineErrorInfo && (
<Box mt={4}>
<ChatErrorCard title={inlineErrorInfo.title} message={inlineErrorInfo.message} />
</Box>
......
......@@ -17,6 +17,7 @@ import {
hasAiInteractiveContent,
hasAiProcessingContent
} from './AIChatBubble/utils';
import { getChatItemRenderKey } from '../utils/recordGroups';
export type ChatRecordsListProps = {
records: ChatSiteItemType[];
......@@ -160,7 +161,7 @@ const ChatRecordsList = ({
previousRecord?.obj === ChatRoleEnum.Human ? onRetry(previousRecord.dataId) : undefined;
return (
<Box key={item.dataId}>
<Box key={getChatItemRenderKey(item)}>
{item.collapseTop && (
<DeletedItemsCollapse
count={item.collapseTop.count}
......
......@@ -15,6 +15,9 @@ type UseChatRecordActionsProps = {
onDeleteChatItem?: (contentId: string, delFile?: boolean) => Promise<void>;
};
const uniqueDataIds = (dataIds: Array<string | undefined>) =>
Array.from(new Set(dataIds.filter((dataId): dataId is string => !!dataId)));
/**
* 管理 ChatBox 中“记录级动作”的副作用和本地 records 更新。
*
......@@ -47,20 +50,23 @@ export const useChatRecordActions = ({
const outLinkAuthData = useContextSelector(WorkflowRuntimeContext, (v) => v.outLinkAuthData);
/**
* 删除一服务端聊天记录。
* 删除一服务端聊天记录。
*
* `delFile` 默认是 true,表示删除消息时一起删除关联文件;重试流程会传 false,
* 因为同一轮历史可能需要继续复用原文件输入,不能在删除旧记录时把文件也删掉。
*/
const onDelMessage = useMemoizedFn((contentId: string, delFile = true) => {
const onDelMessages = useMemoizedFn((contentIds: string[], delFile = true) => {
const targetContentIds = uniqueDataIds(contentIds);
if (targetContentIds.length === 0) return Promise.resolve();
if (onDeleteChatItem) {
return onDeleteChatItem(contentId, delFile);
return Promise.all(targetContentIds.map((contentId) => onDeleteChatItem(contentId, delFile)));
}
return delChatRecordById({
appId,
chatId,
contentId,
contentIds: targetContentIds,
delFile,
...outLinkAuthData
});
......@@ -85,12 +91,9 @@ export const useChatRecordActions = ({
const delHistory = chatRecords.slice(index);
try {
await Promise.all(
delHistory.map((item) => {
if (item.dataId) {
return onDelMessage(item.dataId, false);
}
})
await onDelMessages(
delHistory.map((item) => item.dataId),
false
);
setChatRecords((state) => (index === 0 ? [] : state.slice(0, index)));
......@@ -126,12 +129,9 @@ export const useChatRecordActions = ({
try {
if (index < 0 || delHistory[0]?.obj !== ChatRoleEnum.Human) return;
await Promise.all(
delHistory.map((item) => {
if (item.dataId) {
return onDelMessage(item.dataId, false);
}
})
await onDelMessages(
delHistory.map((item) => item.dataId),
false
);
setChatRecords((state) => (index === 0 ? [] : state.slice(0, index)));
......@@ -162,18 +162,22 @@ export const useChatRecordActions = ({
return () => {
setChatRecords((state) => {
let aiIndex = -1;
const deletedDataIds: string[] = [];
return state.filter((chat, i) => {
const nextState = state.filter((chat, i) => {
if (chat.dataId === dataId) {
aiIndex = i + 1;
onDelMessage(dataId);
deletedDataIds.push(dataId);
return false;
} else if (aiIndex === i && chat.obj === ChatRoleEnum.AI && chat.dataId) {
onDelMessage(chat.dataId);
deletedDataIds.push(chat.dataId);
return false;
}
return true;
});
void onDelMessages(deletedDataIds);
return nextState;
});
};
});
......
......@@ -9,7 +9,7 @@ import {
ChatRoleEnum,
ChatStatusEnum
} from '@fastgpt/global/core/chat/constants';
import { mergeChatResponseData } from '@fastgpt/global/core/chat/utils';
import { mergeNodeResponseDataByIdAndParent } from '@fastgpt/global/core/chat/utils/mergeNode';
import { streamResumeFetch, type ResumeStreamErrorType } from '@/web/common/api/fetch';
import { ChatItemContext } from '@/web/core/chat/context/chatItemContext';
import { ChatRecordContext } from '@/web/core/chat/context/chatRecordContext';
......@@ -276,7 +276,7 @@ export const useChatResume = ({
...item,
status: ChatStatusEnum.finish,
time: new Date(),
responseData: mergeChatResponseData(item.responseData || [])
responseData: mergeNodeResponseDataByIdAndParent(item.responseData || [])
};
});
......
import { useEffect, useMemo, useRef, useState, type MutableRefObject } from 'react';
import { useMemoizedFn, useThrottleFn } from 'ahooks';
import {
isChatScrollAtBottom,
shouldFollowGeneratingScroll,
shouldShowChatScrollToBottomButton
} from '../utils/scrollUtils';
import { isChatScrollAtBottom, shouldShowChatScrollToBottomButton } from '../utils/scrollUtils';
/**
* 管理 ChatBox 的滚动容器和生成中跟随底部逻辑。
......@@ -16,15 +12,17 @@ import {
* 输出约定:
* - `ScrollContainerRef` 绑定到聊天记录滚动容器。
* - `scrollToBottom` 用于明确要求滚到底部,支持 smooth/auto 和延迟。
* - `generatingScroll` 用于流式生成中“条件跟随底部”,避免打断用户查看历史。
* - `generatingScroll` 用于流式生成中根据用户意图跟随底部,避免打断用户查看历史。
* - `isScrollAtBottom` 只描述当前滚动位置。
* - `isScrollToBottomButtonVisible` 由安全距离阈值控制,避免贴近底部时按钮闪烁。
* - `isScrollToBottomButtonVisible` 只在用户主动离开底部后展示,避免流式内容追加或图片
* 加载造成瞬时距离变化时按钮闪烁。
*
* 关键边界:
* - `scrollToBottom` 保留 DOM 未挂载时的延迟重试,因为 ChatBox 有动态加载记录、
* home/chat 分支切换和恢复生成占位消息,调用时机可能早于滚动容器真实出现。
* - 用户主动向上滚动后,当前轮生成不再自动吸附底部;滚回底部或显式点击回到底部后恢复。
* - `generatingScroll` 复用 `shouldFollowGeneratingScroll`,只有仍允许跟随且用户接近底部时才跟随。
* - 内容高度变化时只看 `shouldFollowGeneratingRef`,不再用底部距离反推吸底意图,避免图片加载等
* 异步撑高内容时丢失吸底。
*/
export const useChatScroll = () => {
const scrollContainerRef = useRef<HTMLDivElement | null>(null);
......@@ -86,12 +84,37 @@ export const useChatScroll = () => {
setIsScrollAtBottom(nextIsScrollAtBottom);
setIsScrollToBottomButtonVisible(
!nextIsScrollAtBottom && shouldShowChatScrollToBottomButton(scrollState)
!nextIsScrollAtBottom &&
shouldShowChatScrollToBottomButton({
...scrollState,
userHasLeftBottom: !shouldFollowGeneratingRef.current
})
);
}
);
/**
* 内容尺寸变化后的吸底补偿。
*
* 图片加载、代码块渲染和 Markdown 重排都会在用户没有滚动的情况下改变 scrollHeight;
* 只要用户意图仍是吸底,就直接补滚到底部。若用户已主动上滚,`shouldFollowGeneratingRef`
* 为 false,不会打断阅读历史。
*/
const syncAfterContentChange = useMemoizedFn(() => {
syncScrollAtBottom();
if (!shouldFollowGeneratingRef.current) return;
const container = scrollContainerRef.current;
if (!container) return;
container.scrollTo({
top: container.scrollHeight,
behavior: 'auto'
});
syncScrollAtBottom();
});
/**
* 滚动到底部。
*
* `delay` 用于等待 React 渲染新消息或新高度;内部 `runScroll` 的二次重试用于等待
......@@ -120,20 +143,9 @@ export const useChatScroll = () => {
const { run: generatingScroll } = useThrottleFn(
(force?: boolean) => {
if (!scrollContainerRef.current) return;
if (!shouldFollowGeneratingRef.current) return;
// 流式响应会高频触发,先判断用户是否仍在底部附近,再决定是否滚动。
// `force` 只绕过距离阈值,不绕过用户主动向上滚动后的暂停跟随状态。
const isBottom = shouldFollowGeneratingScroll({
scrollTop: scrollContainerRef.current.scrollTop,
clientHeight: scrollContainerRef.current.clientHeight,
scrollHeight: scrollContainerRef.current.scrollHeight,
force
});
if (!shouldFollowGeneratingRef.current && !force) return;
if (isBottom) {
scrollToBottom('auto');
}
scrollToBottom('auto');
},
{
wait: 100
......@@ -149,7 +161,7 @@ export const useChatScroll = () => {
syncScrollAtBottom();
const handleScroll = () => syncScrollAtBottom({ fromScroll: true });
const handleResize = () => syncScrollAtBottom();
const handleResize = () => syncAfterContentChange();
container.addEventListener('scroll', handleScroll, { passive: true });
window.addEventListener('resize', handleResize);
......@@ -158,7 +170,7 @@ export const useChatScroll = () => {
typeof ResizeObserver === 'undefined'
? undefined
: new ResizeObserver(() => {
syncScrollAtBottom();
syncAfterContentChange();
});
resizeObserver?.observe(container);
......@@ -166,7 +178,7 @@ export const useChatScroll = () => {
typeof MutationObserver === 'undefined'
? undefined
: new MutationObserver(() => {
syncScrollAtBottom();
syncAfterContentChange();
});
mutationObserver?.observe(container, {
childList: true,
......@@ -180,7 +192,7 @@ export const useChatScroll = () => {
resizeObserver?.disconnect();
mutationObserver?.disconnect();
};
}, [scrollContainer, syncScrollAtBottom]);
}, [scrollContainer, syncAfterContentChange, syncScrollAtBottom]);
return {
ScrollContainerRef,
......
/**
* 控制聊天气泡底部错误卡片的展示时机。
* 正在生成的最后一条消息可能先收到节点错误再继续产出内容,生成中先隐藏,结束后再展示最终错误。
*/
export const shouldShowChatItemInlineError = ({
hasInlineError,
isChatting,
isLastChild
}: {
hasInlineError: boolean;
isChatting: boolean;
isLastChild: boolean;
}) => {
return hasInlineError && !(isChatting && isLastChild);
};
......@@ -2,6 +2,15 @@ import { ChatTypeEnum } from '../constants';
import type { ChatSiteItemType } from '../type';
/**
* 生成聊天记录列表使用的 React key。
*
* 同一轮交互下 human/AI 消息可能复用 `dataId`,列表 key 需要带上 `obj` 才能避免
* React 把不同角色的消息当成同一个元素复用。
*/
export const getChatItemRenderKey = (item: Pick<ChatSiteItemType, 'obj' | 'dataId'>) =>
`${item.obj}-${item.dataId}`;
/**
* 为 log 模式的连续删除消息补充折叠元信息。
*
* 输入输出约定:
......
......@@ -4,7 +4,10 @@ import type {
AIChatItemValueItemType,
ChatHistoryItemResType
} from '@fastgpt/global/core/chat/type';
import { appendNodeResponseByParent, mergeChatResponseData } from '@fastgpt/global/core/chat/utils';
import {
appendNodeResponseByParent,
mergeNodeResponseDataByIdAndParent
} from '@fastgpt/global/core/chat/utils/mergeNode';
import { extractDeepestInteractive } from '@fastgpt/global/core/workflow/runtime/utils';
import type { WorkflowInteractiveResponseType } from '@fastgpt/global/core/workflow/template/system/interactive/type';
import type { ChatSiteItemType } from '../type';
......@@ -150,7 +153,7 @@ export const mergeResumeCompletedChatRecords = ({
const mergedResponseData =
shouldMergeResponseData && resumedResponseData?.length
? mergeChatResponseData(
? mergeNodeResponseDataByIdAndParent(
resumedResponseData.reduce<ChatHistoryItemResType[]>(
(responses, resumedItem) => appendNodeResponseByParent(responses, resumedItem),
item.responseData || []
......
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