Commit c063d733 by Archer Committed by GitHub

Perf code (#6978)

* doc

* perf: count token

* doc
parent e927b94c
......@@ -9,7 +9,7 @@ FastGPT 文档采用双文件 i18n 方案,中文为源语言,英文为目标
## 文件结构
文档位于 `document/content/docs/` 目录下
文档位于 `document/content/` 目录下,可能分布在 `docs/``self-host/` 等子目录
- 内容文件:`{name}.mdx`(中文) → `{name}.en.mdx`(英文)
- 导航文件:`meta.json`(中文) → `meta.en.json`(英文)
......@@ -27,6 +27,8 @@ FastGPT 文档采用双文件 i18n 方案,中文为源语言,英文为目标
**手动指定**:用户直接给出文件路径或目录。
无论使用哪种方式,只要翻译范围包含 `.mdx` 内容文件,都必须把同目录的 `meta.json` / `meta.en.json` 纳入检查范围,避免新增英文内容后导航缺失。
### 2. 翻译内容文件(.mdx → .en.mdx)
对每个中文 `.mdx` 文件,生成或更新对应的 `.en.mdx` 文件。
......@@ -47,9 +49,16 @@ FastGPT 文档采用双文件 i18n 方案,中文为源语言,英文为目标
- 表格中的文字内容
- 代码块中的中文注释
### 3. 翻译导航文件(meta.json → meta.en.json)
### 3. 同步导航文件(meta.json → meta.en.json)
对每个中文 `meta.json`,生成或更新对应的 `meta.en.json`。翻译 `.mdx` 文件时,也要同步检查其同目录导航文件:
对每个中文 `meta.json`,生成或更新对应的 `meta.en.json`
- 如果同目录存在 `meta.json`,检查本次翻译的中文文件 basename(如 `41503.mdx``41503`)是否在 `pages` 数组中。
- 如果中文文件是新增页面且 `meta.json` 未引用,应按目录内现有排序和上下文更新 `meta.json`;无法可靠判断位置时,先询问用户,不要只创建 `.en.mdx` 后忽略导航。
- `meta.en.json``pages` 必须与 `meta.json` 保持一致,只翻译 `title``description` 和分隔符字符串,不要翻译或删除页面文件名引用。
- 如果 `meta.en.json` 缺失,必须创建;如果已存在但 `pages` 落后于 `meta.json`,必须同步。
- 如果用户只指定了 `meta.json`,仍按导航文件翻译规则更新 `meta.en.json`,并检查 `pages` 中引用的中文 `.mdx` 是否有对应 `.en.mdx`
- 修改导航 JSON 时优先使用结构化 JSON 方式或严格保持原文件格式,避免手改导致尾逗号、缩进漂移或 `pages` 顺序错误。
**需要翻译的字段**`title``description`、分隔符字符串(如 `"---入门---"``"---Getting Started---"`
......@@ -58,7 +67,9 @@ FastGPT 文档采用双文件 i18n 方案,中文为源语言,英文为目标
### 4. 翻译完成后
- 列出所有已翻译的文件
- 列出所有已更新或确认无需更新的 `meta.json` / `meta.en.json`
- 如果发现中文文件有对应英文文件缺失的情况,提醒用户
- 如果发现中文 `.mdx` 未被同目录 `meta.json` 引用,且本次未能安全更新导航,明确提醒用户
## 翻译原则
......
......@@ -22,7 +22,6 @@ The example below includes system parameters and model configurations:
"vectorMaxProcess": 15, // Vector processing thread count
"qaMaxProcess": 15, // Q&A splitting thread count
"vlmMaxProcess": 15, // Vision-language model max processing threads
"tokenWorkers": 50, // Token calculation worker count — keeps memory occupied, don't set too high
"hnswEfSearch": 100, // Vector search parameter (PG and OB only). Higher = more accurate but slower. 100 gives 99%+ accuracy.
"customPdfParse": {
// Added in v4.9.0
......
......@@ -13,7 +13,7 @@ description: FastGPT config.json 文件配置
## 4.8.20+ 版本新配置文件示例
> 从4.8.20版本开始,模型在页面中进行配置。
> 从 4.8.20 版本开始,模型在页面中进行配置。
```json
{
......@@ -24,7 +24,6 @@ description: FastGPT config.json 文件配置
"vectorMaxProcess": 15, // 向量处理线程数量
"qaMaxProcess": 15, // 问答拆分线程数量
"vlmMaxProcess": 15, // 图片理解模型最大处理进程
"tokenWorkers": 50, // Token 计算线程保持数,会持续占用内存,不能设置太大。
"hnswEfSearch": 100, // 向量搜索参数,仅对 PG 和 OB 生效。越大,搜索越精确,但是速度越慢。设置为100,有99%+精度。
"customPdfParse": {
// 4.9.0 新增配置
......@@ -45,16 +44,16 @@ description: FastGPT config.json 文件配置
#### 1. 申请 Sealos AI proxy API Key
[点击打开 Sealos Pdf parser 官网](https://hzh.sealos.run/?uid=fnWRt09fZP&openapp=system-aiproxy),并进行对应 API Key 的申请。
[点击打开 Sealos PDF parser 官网](https://hzh.sealos.run/?uid=fnWRt09fZP&openapp=system-aiproxy),并进行对应 API Key 的申请。
#### 2. 修改 FastGPT 配置文件
`systemEnv.customPdfParse.url`填写成`https://aiproxy.hzh.sealos.run/v1/parse/pdf?model=parse-pdf`
`systemEnv.customPdfParse.key`填写成在 Sealos AI proxy 中申请的 API Key。
`systemEnv.customPdfParse.url` 填写成 `https://aiproxy.hzh.sealos.run/v1/parse/pdf?model=parse-pdf`
`systemEnv.customPdfParse.key` 填写成在 Sealos AI proxy 中申请的 API Key。
### 使用 Doc2x 解析 PDF 文件
`Doc2x`是一个国内提供专业 PDF 解析。
`Doc2x` 是一个国内提供专业 PDF 解析。
#### 1. 申请 Doc2x 服务
......@@ -68,7 +67,7 @@ description: FastGPT config.json 文件配置
#### 3. 开始使用
在知识库导入数据或应用文件上传配置中,可以勾选`PDF 增强解析`,则在对 PDF 解析时候,会使用 Doc2x 服务进行解析。
在知识库导入数据或应用文件上传配置中,可以勾选 `PDF 增强解析`,则在对 PDF 解析时候,会使用 Doc2x 服务进行解析。
### 使用 Marker 解析 PDF 文件
......
---
title: 'V4.15.0-beta2'
description: 'FastGPT V4.15.0-beta2 Release Notes'
---
## 📦 Upgrade Guide
### Image Changes
- Update the fastgpt-app (FastGPT main service) image tag to v4.15.0-beta2
- Update the fastgpt-pro (FastGPT commercial edition) image tag to v4.15.0-beta2
If you use OpenSandbox, update the following images:
- Update the fastgpt-agent-sandbox image tag to v0.2.0
- Update the fastgpt-agent-volume-manager image tag to v0.2.0
### Environment Variable Updates
1. If you use OpenSandbox, `AGENT_SANDBOX_VOLUME_MANAGER_MOUNT_PATH` is no longer effective and can be removed. OpenSandbox now always mounts persistent data to `/workspace`, which affects the old sandbox persistence behavior.
## 🚀 New Features
1. Added Skill editing. Agents can now use Skills. Currently, only static Skills are supported, and reverse calls to system tools are not supported.
2. Reworked the agentV2 loop logic.
3. Knowledge Base search now supports native multimodal embedding models and image-to-image search.
4. Chat API `dataId` validation: `/v1/chat/completions`, `/v2/chat/completions`, and `chatTest` now validate whether the current `dataId` duplicates one in the request or existing records in the current session before running the Workflow. Duplicate values return a business error immediately, preventing invalid data from entering Workflow execution and stream-resume merge logic.
## ⚙️ Improvements
1. Optimized the OTEL log collection format.
2. Disabled invalid connection mode in Workflows.
3. Improved layout adaptation for Workflow nodes with very long names.
4. Improved the Knowledge Base search test interaction.
5. Improved the Knowledge Base data editing modal.
6. Improved the reason hide toggle so reasoning can be hidden in the UI while still being preserved when requesting the LLM.
7. Stream-resume pause experience: after pausing, the client waits for the backend to return the real generation state. If the Workflow has not finished, the input area remains disabled and shows "Stopping", preventing the next round from being sent before the previous one ends.
8. Faster recovery for abnormally interrupted sessions: after a service crash or restart, Redis stream activity detection (about two minutes without a heartbeat) is used to correct stuck "generating" sessions to completed sooner. The 30-minute MongoDB fallback is still retained, and short Redis outages will not incorrectly update sessions that are still generating.
9. Remember the most recent chat when switching apps: when switching apps in the same browser, the last opened `chatId` is restored per app instead of sharing a single global session id.
10. Optimized response detail display: in the full response modal, file fields from form input nodes are displayed as file lists instead of raw JSON text.
11. Optimized `chat2messages` adaptation to avoid standalone reason output.
## 🐛 Bug Fixes
1. Fixed abnormal default values in Workflow single-node debugging.
2. Fixed abnormal `defaultConfig` override behavior in model configuration.
3. Clear the local chat cache when switching teams.
4. Conversation stream resume:
- Submitted form input values, including `fileSelect` file lists, are correctly restored into interactive nodes after refresh or reconnect resume. Empty forms and disappearing files no longer occur.
- Loaded AI output and node responses are preserved when automatic resume starts. When completed records overwrite local state, restored interactive form values and flow node responses are no longer lost.
- Expired unsubmitted interactions are no longer appended again after a form is submitted. During resume, form default values now stay in sync with `formInputResult`.
- After starting a new conversation, the temporary sidebar history item prioritizes the title generated from user input. The server-side title overwrites it after being persisted, avoiding a long-running "New Chat" display.
- Fixed an issue where the sidebar or conversation content briefly showed chat records from another app when switching apps.
5. Stop conversation prompt: removed the warning toast shown during stop and replaced it with a status prompt synchronized with the backend generation state.
6. Fixed the v1/completions API where `quoteList` in `nodeResponse` did not return `q` and `a`.
## 🛠️ Code Improvements
1. Split AI request logic and Workflow run detail code.
2. Updated billing logic for user-defined API keys.
3. Added design documentation and unit tests for stream-resume-related modules, including stop state, stale cleanup, history title, `dataId` validation, and form restoration.
4. Changed the volume manager runtime from Bun to Node.js.
---
title: 'V4.15.0-beta2(进行中)'
title: 'V4.15.0-beta2'
description: 'FastGPT V4.15.0-beta2 更新说明'
---
## 升级指南
## 📦 升级指南
### 镜像变更
......@@ -54,7 +54,7 @@ description: 'FastGPT V4.15.0-beta2 更新说明'
5. 停止对话提示:移除停止时的 warning toast,改为与后端生成态同步的状态提示。
6. v1/completions 接口,nodeResponse 中,quoteList 未返回 `q` , `a`。
## 代码优化
## 🛠️ 代码优化
1. 拆分 AI request、工作流运行详情代码。
2. 用户自定义密钥计费逻辑。
......
---
title: 'V4.15.0-beta3 (In Progress)'
description: 'FastGPT V4.15.0-beta3 Release Notes'
---
## 🚀 New Features
## ⚙️ Improvements
## 🐛 Bug Fixes
## 🛠️ Code Improvements
1. Updated the token calculation dependency to improve performance.
---
title: 'V4.15.0-beta3(进行中)'
description: 'FastGPT V4.15.0-beta3 更新说明'
---
## 🚀 新增内容
## ⚙️ 优化
## 🐛 修复
## 🛠️ 代码优化
1. 调整 token 计算依赖,提高性能。
{
"title": "4.15.x",
"description": "",
"pages": ["41502", "4150"]
"pages": ["41503", "41502", "4150"]
}
{
"title": "4.15.x",
"description": "",
"pages": ["41502", "4150"]
"pages": ["41503", "41502", "4150"]
}
......@@ -140,6 +140,8 @@ description: FastGPT Toc
- [/en/self-host/upgrading/4-14/4148](/en/self-host/upgrading/4-14/4148)
- [/en/self-host/upgrading/4-14/41481](/en/self-host/upgrading/4-14/41481)
- [/en/self-host/upgrading/4-14/4149](/en/self-host/upgrading/4-14/4149)
- [/en/self-host/upgrading/4-15/41502](/en/self-host/upgrading/4-15/41502)
- [/en/self-host/upgrading/4-15/41503](/en/self-host/upgrading/4-15/41503)
- [/en/self-host/upgrading/outdated/40](/en/self-host/upgrading/outdated/40)
- [/en/self-host/upgrading/outdated/41](/en/self-host/upgrading/outdated/41)
- [/en/self-host/upgrading/outdated/4100](/en/self-host/upgrading/outdated/4100)
......
......@@ -144,6 +144,7 @@ description: FastGPT 文档目录
- [/self-host/upgrading/4-14/4149](/self-host/upgrading/4-14/4149)
- [/self-host/upgrading/4-15/4150](/self-host/upgrading/4-15/4150)
- [/self-host/upgrading/4-15/41502](/self-host/upgrading/4-15/41502)
- [/self-host/upgrading/4-15/41503](/self-host/upgrading/4-15/41503)
- [/self-host/upgrading/outdated/40](/self-host/upgrading/outdated/40)
- [/self-host/upgrading/outdated/41](/self-host/upgrading/outdated/41)
- [/self-host/upgrading/outdated/4100](/self-host/upgrading/outdated/4100)
......
......@@ -159,8 +159,8 @@
"content/openapi/share.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/config/env.en.mdx": "2026-05-23T22:47:02+08:00",
"content/self-host/config/env.mdx": "2026-05-23T22:47:02+08:00",
"content/self-host/config/json.en.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/config/json.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/config/json.en.mdx": "2026-05-25T10:13:06+08:00",
"content/self-host/config/json.mdx": "2026-05-25T10:13:06+08:00",
"content/self-host/config/model/intro.en.mdx": "2026-05-07T15:06:40+08:00",
"content/self-host/config/model/intro.mdx": "2026-05-07T15:06:40+08:00",
"content/self-host/config/model/minimax.en.mdx": "2026-04-26T21:08:47+08:00",
......@@ -278,7 +278,7 @@
"content/self-host/upgrading/4-14/4149.en.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/4-14/4149.mdx": "2026-04-26T21:08:47+08:00",
"content/self-host/upgrading/4-15/4150.mdx": "2026-05-20T17:52:26+08:00",
"content/self-host/upgrading/4-15/41502.mdx": "2026-05-24T17:01:11+08:00",
"content/self-host/upgrading/4-15/41502.mdx": "2026-05-24T18:19:34+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",
......
......@@ -9,7 +9,6 @@
"vectorMaxProcess": 10, // 向量处理线程数量
"qaMaxProcess": 10, // 问答拆分线程数量
"vlmMaxProcess": 10, // 图片理解模型最大处理进程
"tokenWorkers": 30, // Token 计算线程保持数,会持续占用内存,不能设置太大。
"hnswEfSearch": 100, // 向量搜索参数,仅对 PG OB 生效。越大,搜索越精确,但是速度越慢。设置为100,有99%+精度。
"hnswMaxScanTuples": 100000, // 向量搜索最大扫描数据量,仅对 PG生效。
"customPdfParse": {
......
......@@ -150,7 +150,6 @@ export type FastGPTFeConfigsType = {
export type SystemEnvType = {
openapiPrefix?: string;
tokenWorkers: number; // token count max worker (min 10, max 1000)
datasetParseMaxProcess: number;
vectorMaxProcess: number;
......
......@@ -7,62 +7,121 @@ import {
import { chats2GPTMessages } from '@fastgpt/global/core/chat/adapt';
import { type ChatItemMiniType } from '@fastgpt/global/core/chat/type';
import { WorkerNameEnum, getWorkerController } from '../../../worker/utils';
import { getTokenWorkerCount } from '../../../worker/tokenWorkerConfig';
import type { ChatCompletionRequestMessageRoleEnum } from '@fastgpt/global/core/ai/constants';
import { getLogger, LogCategories } from '../../logger';
const logger = getLogger(LogCategories.MODULE.AI.LLM);
export const countGptMessagesTokens = async (
messages: ChatCompletionMessageParam[],
tools?: ChatCompletionTool[],
functionCall?: ChatCompletionCreateParams.Function[]
) => {
try {
const workerController = getWorkerController<
{
messages: ChatCompletionMessageParam[];
tools?: ChatCompletionTool[];
functionCall?: ChatCompletionCreateParams.Function[];
},
number
>({
name: WorkerNameEnum.countGptMessagesTokens,
maxReservedThreads: global.systemEnv?.tokenWorkers || 30
});
export type CountGptMessagesTokensParams = {
messages: ChatCompletionMessageParam[];
tools?: ChatCompletionTool[];
functionCall?: ChatCompletionCreateParams.Function[];
};
type CountGptMessagesTokensWorkerPayload = {
messages?: ChatCompletionMessageParam[];
messageGroups?: ChatCompletionMessageParam[][];
prompts?: (string | null | undefined)[];
tools?: ChatCompletionTool[];
functionCall?: ChatCompletionCreateParams.Function[];
};
const total = await workerController.run({ messages, tools, functionCall });
/**
* 获取 token 计数 worker 池。
*
* 主进程不直接 import tokenizer,避免把 o200k_base 编码表加载到 API 进程常驻内存;
* worker 数量由 getTokenWorkerCount 统一限制,和启动预热逻辑保持一致。
*/
const getTokenCountWorkerController = <Response = number>() =>
getWorkerController<CountGptMessagesTokensWorkerPayload, Response>({
name: WorkerNameEnum.countGptMessagesTokens,
maxReservedThreads: getTokenWorkerCount()
});
return total;
/**
* 统一封装 token worker 调用,保留失败日志的模块上下文。
*
* 这里不做主线程本地 fallback:fallback 会重新加载 tokenizer 到主进程,抵消 worker
* 隔离内存的收益;失败时直接抛出,让上层按正常错误链路处理。
*/
const runTokenCountWorker = async <Response>(payload: CountGptMessagesTokensWorkerPayload) => {
try {
const workerController = getTokenCountWorkerController<Response>();
return await workerController.run(payload);
} catch (error) {
logger.error('Token count worker failed, using fallback', { error });
const total = messages.reduce((sum, item) => {
if (item.content) {
return sum + item.content.length * 0.5;
}
return sum;
}, 0);
return total;
logger.error('Token count worker failed', { error });
throw error;
}
};
/**
* 统计 Chat messages token 数。
*
* 这是业务侧的统一入口,内部固定走 token worker 和 o200k_base 编码;该值用于上下文预算
* 和供应商未返回 usage 时的兜底统计,不能替代供应商真实 usage。
*/
export const countGptMessagesTokens = async ({
messages,
tools,
functionCall
}: CountGptMessagesTokensParams) => {
return runTokenCountWorker<number>({ messages, tools, functionCall });
};
/**
* 批量统计多组 Chat messages token。
*
* 用于上下文裁剪等热路径,避免每一轮对话都单独 postMessage 到 worker。
*/
export const countGptMessagesTokensBatch = async (
messageGroups: ChatCompletionMessageParam[][]
) => {
const totals = await runTokenCountWorker<number[]>({ messageGroups });
if (totals.length !== messageGroups.length) {
throw new Error('Token count worker returned mismatched message group result length');
}
return totals;
};
export const countMessagesTokens = (messages: ChatItemMiniType[]) => {
const adaptMessages = chats2GPTMessages({ messages, reserveId: true });
return countGptMessagesTokens(adaptMessages);
return countGptMessagesTokens({ messages: adaptMessages });
};
/* count one prompt tokens */
/**
* 统计单段普通 prompt token。
*
* 历史调用方会传入空 role,把 prompt 包装成最小 chat message;该兼容行为由 worker 内部
* 处理,避免纯文本 prompt 被额外加上 chat role 固定开销。
*/
export const countPromptTokens = async (
prompt: string | ChatCompletionContentPart[] | null | undefined = '',
role: '' | `${ChatCompletionRequestMessageRoleEnum}` = ''
) => {
const total = await countGptMessagesTokens([
{
//@ts-ignore
role,
content: prompt
}
]);
const total = await countGptMessagesTokens({
messages: [
{
//@ts-ignore
role,
content: prompt
}
]
});
return total;
};
/**
* 批量统计普通 prompt token,主要用于知识库召回和 embedding/rerank 兜底计量。
*/
export const countPromptTokensBatch = async (prompts: (string | null | undefined)[]) => {
const totals = await runTokenCountWorker<number[]>({ prompts });
if (totals.length !== prompts.length) {
throw new Error('Token count worker returned mismatched prompt result length');
}
return totals;
};
......@@ -547,7 +547,9 @@ export const runAgentLoop = async <TChildrenResponse = unknown>({
while (toolIndex < toolCalls.length) {
const currentTool = toolCalls[toolIndex];
const currentMessagesTokens = await countGptMessagesTokens(requestMessages);
const currentMessagesTokens = await countGptMessagesTokens({
messages: requestMessages
});
// 只能串行的工具
if (!canBatchTool(currentTool)) {
......
......@@ -98,7 +98,9 @@ export const compressRequestMessages = async ({
// 触发阈值按完整请求上下文判断;压缩内容仍只包含非 system/developer 历史。
// system/developer 虽然不参与 checkpoint 压缩,但会真实占用模型上下文。
const messageTokens = await countGptMessagesTokens(messages);
const messageTokens = await countGptMessagesTokens({
messages
});
const thresholds = calculateCompressionThresholds(model.maxContext).messages;
if (messageTokens < thresholds.threshold) {
......@@ -359,10 +361,14 @@ export const compressLargeContent = async ({
let merged = compressedChunks.join('\n\n');
// LLM 输出长度不可控,合并后仍需做一次真实 token 校验。
const finalTokens = await countGptMessagesTokens([{ role: 'user', content: merged }]);
const finalTokens = await countGptMessagesTokens({
messages: [{ role: 'user', content: merged }]
});
logger.info('LLM chunk compression Completed', {
originalTokens: await countGptMessagesTokens([{ role: 'user', content: content }]),
originalTokens: await countGptMessagesTokens({
messages: [{ role: 'user', content: content }]
}),
finalTokens,
compressedTokenLimit,
success: finalTokens <= compressedTokenLimit
......
......@@ -158,12 +158,15 @@ export const createLLMResponse = async <T extends ChatCompletionCreateParams>(
// 但类型层 InferCompletionsBody 落到了 SDK 形态(messages/tools 是 v6 后的 union)
const inputTokens =
usage?.prompt_tokens ||
(await countGptMessagesTokens(
requestBody.messages as ChatCompletionMessageParam[],
requestBody.tools as ChatCompletionTool[] | undefined
));
(await countGptMessagesTokens({
messages: requestBody.messages as ChatCompletionMessageParam[],
tools: requestBody.tools as ChatCompletionTool[] | undefined
}));
const outputTokens =
usage?.completion_tokens || (await countGptMessagesTokens([assistantMessage]));
usage?.completion_tokens ||
(await countGptMessagesTokens({
messages: [assistantMessage]
}));
/**
* 空响应不一定是请求异常,可能是模型返回 stop 但没有内容。
......
import { countGptMessagesTokens } from '../../../common/string/tiktoken/index';
import {
countGptMessagesTokens,
countGptMessagesTokensBatch
} from '../../../common/string/tiktoken/index';
import type {
ChatCompletionAssistantMessageParam,
ChatCompletionContentPart,
......@@ -61,7 +64,9 @@ export const filterGPTMessageByMaxContext = async ({
}
// reduce token of systemPrompt
maxContext -= await countGptMessagesTokens([...systemPrompts, ...leadingContextCheckpoints]);
maxContext -= await countGptMessagesTokens({
messages: [...systemPrompts, ...leadingContextCheckpoints]
});
/* 截取时候保证一轮内容的完整性
1. user - assistant - user
......@@ -70,11 +75,10 @@ export const filterGPTMessageByMaxContext = async ({
3. user - assistant - tool - assistant - tool
4. user - assistant - assistant - tool - tool
*/
// Save the last chat prompt(question)
let chats: ChatCompletionMessageParam[] = [];
const messageGroups: ChatCompletionMessageParam[][] = [];
let tmpChats: ChatCompletionMessageParam[] = [];
// 从后往前截取对话内容, 每次到 user 则认为是一组完整信息
// 从后往前分组,每次到 user 则认为是一组完整信息;后续批量统计可减少 worker 往返。
while (chatPrompts.length > 0) {
const lastMessage = chatPrompts.pop();
if (!lastMessage) {
......@@ -83,20 +87,27 @@ export const filterGPTMessageByMaxContext = async ({
// 遇到 user,说明到了一轮完整信息,可以开始判断是否需要保留
if (lastMessage.role === ChatCompletionRequestMessageRoleEnum.User) {
const tokens = await countGptMessagesTokens([lastMessage, ...tmpChats]);
maxContext -= tokens;
// 该轮信息整体 tokens 超出范围,这段数据不要了。但是至少保证一组。
if (maxContext < 0 && chats.length > 0) {
break;
}
chats = [lastMessage, ...tmpChats].concat(chats);
messageGroups.unshift([lastMessage, ...tmpChats]);
tmpChats = [];
} else {
tmpChats.unshift(lastMessage);
}
}
const reversedMessageGroups = [...messageGroups].reverse();
const groupTokens = (await countGptMessagesTokensBatch(reversedMessageGroups)).reverse();
const chats: ChatCompletionMessageParam[] = [];
for (let i = messageGroups.length - 1; i >= 0; i--) {
maxContext -= groupTokens[i] || 0;
// 该轮信息整体 tokens 超出范围,这段数据不要了。但是至少保证一组。
if (maxContext < 0 && chats.length > 0) {
break;
}
chats.unshift(...messageGroups[i]);
}
return [...systemPrompts, ...leadingContextCheckpoints, ...chats];
};
......
......@@ -149,7 +149,7 @@ export async function reRankRecall({
results,
inputTokens:
data?.meta?.tokens?.input_tokens ||
(await countPromptTokens(documentsTextArray.join('\n') + query, ''))
(await countPromptTokens(documentsTextArray.join('\n') + query))
};
})
.catch((err) => {
......
import { DatasetSearchModeEnum } from '@fastgpt/global/core/dataset/constants';
import type { SearchDataResponseItemType } from '@fastgpt/global/core/dataset/type';
import { countPromptTokens } from '../../../../common/string/tiktoken/index';
import { countPromptTokensBatch } from '../../../../common/string/tiktoken/index';
import { getLogger, LogCategories } from '../../../../common/logger';
const logger = getLogger(LogCategories.MODULE.DATASET.DATA);
/**
* 根据搜索模式分配每条召回链路的候选数量。
......@@ -37,12 +40,12 @@ export const filterDatasetDataByMaxTokens = async (
data: SearchDataResponseItemType[],
maxTokens: number
) => {
const tokensScoreFilter = await Promise.all(
data.map(async (item) => ({
...item,
tokens: await countPromptTokens(item.q + item.a)
}))
);
const startTime = Date.now();
const tokenList = await countPromptTokensBatch(data.map((item) => item.q + item.a));
const tokensScoreFilter = data.map((item, index) => ({
...item,
tokens: tokenList[index] || 0
}));
const results: SearchDataResponseItemType[] = [];
let totalTokens = 0;
......@@ -57,5 +60,20 @@ export const filterDatasetDataByMaxTokens = async (
}
}
return results.length === 0 ? data.slice(0, 1) : results;
const filteredResults = results.length === 0 ? data.slice(0, 1) : results;
const obj = {
candidateCount: data.length,
resultCount: filteredResults.length,
maxTokens,
totalTokens,
durationMs: Date.now() - startTime
};
if (obj.durationMs > 100) {
logger.warn('Dataset search token filter completed', obj);
} else {
logger.debug('Dataset search token filter completed', obj);
}
return filteredResults;
};
import { type SearchDataResponseItemType } from '@fastgpt/global/core/dataset/type';
import { countPromptTokens } from '../../../common/string/tiktoken/index';
import { countPromptTokensBatch } from '../../../common/string/tiktoken/index';
import type { RuntimeNodeItemType } from '@fastgpt/global/core/workflow/runtime/type';
import { getSystemToolByIdAndVersionId, getSystemTools } from '../../app/tool/controller';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
......@@ -14,10 +14,11 @@ export const filterSearchResultsByMaxChars = async (
) => {
const results: SearchDataResponseItemType[] = [];
let totalTokens = 0;
const itemTokens = await countPromptTokensBatch(list.map((item) => item.q + item.a));
for (let i = 0; i < list.length; i++) {
const item = list[i];
totalTokens += await countPromptTokens(item.q + item.a);
totalTokens += itemTokens[i] || 0;
if (totalTokens > maxTokens + 500) {
break;
}
......
......@@ -39,6 +39,7 @@
"encoding": "^0.1.13",
"file-type": "catalog:",
"form-data": "^4.0.4",
"gpt-tokenizer": "catalog:",
"http-proxy-agent": "^7.0.2",
"https-proxy-agent": "^7.0.6",
"iconv-lite": "^0.6.3",
......@@ -71,7 +72,6 @@
"proxy-addr": "catalog:",
"proxy-agent": "catalog:",
"proxy-from-env": "^1.1.0",
"tiktoken": "1.0.17",
"turndown": "^7.1.2",
"undici": "^7.18.2",
"winston": "^3.17.0",
......@@ -94,4 +94,4 @@
"@types/tunnel": "^0.0.4",
"@types/turndown": "^5.0.4"
}
}
\ No newline at end of file
}
import { describe, expect, it } from 'vitest';
import {
GPT_TOKENIZER_ENCODING,
countGptMessagesTokensInWorker,
countPromptTokensInWorker
} from '@fastgpt/service/worker/countGptMessagesTokens/count';
describe('token counter', () => {
it('should use the modern GPT o200k_base encoding', () => {
const text =
'FastGPT 是一个 AI Agent 构建平台,通过 Flow 提供数据处理、模型调用和可视化工作流编排。';
expect(GPT_TOKENIZER_ENCODING).toBe('o200k_base');
expect(countPromptTokensInWorker(text)).toBe(28);
});
it('should keep prompt-only counts compatible with empty role messages', () => {
const text = 'FastGPT 是一个 AI Agent 构建平台。';
const promptTokens = countPromptTokensInWorker(text);
const messageTokens = countGptMessagesTokensInWorker({
messages: [
{
// countPromptTokens 历史上会把普通 prompt 包成空 role message。
// 该路径不能额外计算 chat message 的 role 固定开销。
role: '' as any,
content: text
}
]
});
expect(messageTokens).toBe(promptTokens);
});
it('should count message content, role overhead, tools and assistant calls', () => {
const tokens = countGptMessagesTokensInWorker({
messages: [
{ role: 'system', content: 'You are helpful.' },
{ role: 'user', content: 'Search FastGPT docs.' },
{
role: 'assistant',
content: '',
tool_calls: [
{
id: 'call_1',
type: 'function',
function: {
name: 'search',
arguments: '{"query":"FastGPT"}'
}
}
]
}
],
tools: [
{
type: 'function',
function: {
name: 'search',
description: 'Search docs',
parameters: {
type: 'object',
properties: {
query: { type: 'string' }
},
required: ['query']
}
}
}
]
});
expect(tokens).toBeGreaterThan(40);
});
});
......@@ -6,6 +6,14 @@ import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
// Mock the AI API client factory so tests don't hit the real network.
// We control the embeddings.create implementation per-test via `mockCreate`.
const mockCreate = vi.fn();
const mockCountPromptTokens = vi.hoisted(() => vi.fn(async (text: string) => text.length));
// getVectors 在缺少 usage 时会回退本地 token 计数;测试里只验证回退路径生效,
// 不启动真实 worker,避免 service 包单测依赖 app/pro 的 worker 构建目录。
vi.mock('@fastgpt/service/common/string/tiktoken/index', () => ({
countPromptTokens: mockCountPromptTokens
}));
vi.mock('@fastgpt/service/core/ai/config', () => ({
getAIApi: () => ({
ai: {
......@@ -370,6 +378,7 @@ describe('getVectors function test', () => {
beforeEach(() => {
mockCreate.mockReset();
mockCountPromptTokens.mockImplementation(async (text: string) => text.length);
});
const buildModel = (overrides: Partial<EmbeddingModelItemType> = {}): EmbeddingModelItemType =>
......
......@@ -99,7 +99,7 @@ describe('compressRequestMessages', () => {
mockDefaultUsagePoints();
countPromptTokensMock.mockResolvedValue(100);
countGptMessagesTokensMock.mockImplementation(
async (messages: ChatCompletionMessageParam[]) => messages.length * 1000
async ({ messages }: { messages: ChatCompletionMessageParam[] }) => messages.length * 1000
);
});
......@@ -261,8 +261,9 @@ describe('compressRequestMessages', () => {
content: 'short user history'
}
];
countGptMessagesTokensMock.mockImplementation(async (input: ChatCompletionMessageParam[]) =>
input === messages ? 4000 : 100
countGptMessagesTokensMock.mockImplementation(
async (input: { messages: ChatCompletionMessageParam[] }) =>
input.messages === messages ? 4000 : 100
);
const result = await compressRequestMessages({
......@@ -271,7 +272,9 @@ describe('compressRequestMessages', () => {
});
expect(createLLMResponseMock).toHaveBeenCalledTimes(1);
expect(countGptMessagesTokensMock).toHaveBeenCalledWith(messages);
expect(countGptMessagesTokensMock).toHaveBeenCalledWith({
messages
});
expect(result.messages).toEqual([
messages[0],
messages[1],
......
......@@ -7,9 +7,15 @@ import { ChatCompletionRequestMessageRoleEnum } from '@fastgpt/global/core/ai/co
import { describe, expect, it, vi, beforeEach } from 'vitest';
// Mock external dependencies
vi.mock('@fastgpt/service/common/string/tiktoken/index', () => ({
countGptMessagesTokens: vi.fn()
}));
vi.mock('@fastgpt/service/common/string/tiktoken/index', () => {
const countGptMessagesTokens = vi.fn();
return {
countGptMessagesTokens,
countGptMessagesTokensBatch: vi.fn((messageGroups: ChatCompletionMessageParam[][]) =>
Promise.all(messageGroups.map((messages) => countGptMessagesTokens({ messages })))
)
};
});
vi.mock('@fastgpt/service/common/file/image/utils', () => ({
getImageBase64: vi.fn()
......@@ -235,9 +241,9 @@ describe('filterGPTMessageByMaxContext function tests', () => {
{ role: ChatCompletionRequestMessageRoleEnum.System, content: 'System' },
{ role: ChatCompletionRequestMessageRoleEnum.User, content: 'current user request' }
]);
expect(mockCountGptMessagesTokens).toHaveBeenNthCalledWith(1, [
{ role: ChatCompletionRequestMessageRoleEnum.System, content: 'System' }
]);
expect(mockCountGptMessagesTokens).toHaveBeenNthCalledWith(1, {
messages: [{ role: ChatCompletionRequestMessageRoleEnum.System, content: 'System' }]
});
});
it('should not preserve malformed hidden checkpoint-like messages as leading checkpoints', async () => {
......
......@@ -13,6 +13,9 @@ const mockMongoDatasetCollectionFind = vi.hoisted(() => vi.fn());
const mockMongoDatasetDataFind = vi.hoisted(() => vi.fn());
const mockMongoDatasetDataTextAggregate = vi.hoisted(() => vi.fn());
const mockGetImageBase64 = vi.hoisted(() => vi.fn());
const mockCountPromptTokensBatch = vi.hoisted(() =>
vi.fn(async (prompts: string[]) => prompts.map((prompt) => prompt.length))
);
const originalMultipleDataToBase64 = serviceEnv.MULTIPLE_DATA_TO_BASE64;
......@@ -39,6 +42,12 @@ vi.mock('@fastgpt/service/common/file/image/utils', () => ({
getImageBase64: mockGetImageBase64
}));
// defaultRecall 的结果过滤只关心 token 数的相对大小,测试里用稳定 mock
// 隔离真实 worker 路径,避免单元测试依赖 app/pro 的 worker 构建产物。
vi.mock('@fastgpt/service/common/string/tiktoken/index', () => ({
countPromptTokensBatch: mockCountPromptTokensBatch
}));
vi.mock('@fastgpt/service/core/dataset/collection/schema', () => ({
DatasetColCollectionName: 'dataset_collections',
MongoDatasetCollection: {
......@@ -69,6 +78,9 @@ describe('default recall dataset search', () => {
beforeEach(() => {
vi.clearAllMocks();
serviceEnv.MULTIPLE_DATA_TO_BASE64 = originalMultipleDataToBase64;
mockCountPromptTokensBatch.mockImplementation(async (prompts: string[]) =>
prompts.map((prompt) => prompt.length)
);
mockGetEmbeddingModel.mockReturnValue({
model: 'mock-embedding-model',
......
......@@ -3,11 +3,11 @@ import type { SearchDataResponseItemType } from '@fastgpt/global/core/dataset/ty
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import { NodeInputKeyEnum } from '@fastgpt/global/core/workflow/constants';
// Mock countPromptTokens (uses worker threads, not available in test)
// Mock token counting (uses worker threads, not available in test)
vi.mock('@fastgpt/service/common/string/tiktoken/index', () => ({
countPromptTokens: vi.fn(async (text: string) => {
countPromptTokensBatch: vi.fn(async (texts: string[]) => {
// Simple approximation: 1 token ≈ 1 char for test purposes
return typeof text === 'string' ? text.length : 0;
return texts.map((text) => (typeof text === 'string' ? text.length : 0));
})
}));
......@@ -21,13 +21,13 @@ import {
filterSearchResultsByMaxChars,
getSystemToolRunTimeNodeFromSystemToolset
} from '@fastgpt/service/core/workflow/utils';
import { countPromptTokens } from '@fastgpt/service/common/string/tiktoken/index';
import { countPromptTokensBatch } from '@fastgpt/service/common/string/tiktoken/index';
import {
getSystemTools,
getSystemToolByIdAndVersionId
} from '@fastgpt/service/core/app/tool/controller';
const mockedCountPromptTokens = vi.mocked(countPromptTokens);
const mockedCountPromptTokensBatch = vi.mocked(countPromptTokensBatch);
const mockedGetSystemTools = vi.mocked(getSystemTools);
const mockedGetSystemToolByIdAndVersionId = vi.mocked(getSystemToolByIdAndVersionId);
......@@ -118,7 +118,7 @@ describe('filterSearchResultsByMaxChars', () => {
it('should concatenate q and a for token counting', async () => {
const list = [makeSearchItem('hello', 'world')];
await filterSearchResultsByMaxChars(list, 100);
expect(mockedCountPromptTokens).toHaveBeenCalledWith('helloworld');
expect(mockedCountPromptTokensBatch.mock.calls[0][0]).toEqual(['helloworld']);
});
it('should handle items with undefined a field', async () => {
......@@ -127,7 +127,7 @@ describe('filterSearchResultsByMaxChars', () => {
const list = [item];
await filterSearchResultsByMaxChars(list, 100);
// q + a = "hello" + undefined = "helloundefined"
expect(mockedCountPromptTokens).toHaveBeenCalledWith('helloundefined');
expect(mockedCountPromptTokensBatch.mock.calls[0][0]).toEqual(['helloundefined']);
});
});
......
import { beforeEach, describe, expect, it, vi } from 'vitest';
const osMock = vi.hoisted(() => ({
availableParallelism: vi.fn(),
cpus: vi.fn()
}));
vi.mock('os', () => ({
availableParallelism: osMock.availableParallelism,
cpus: osMock.cpus
}));
import { getTokenWorkerCount } from '@fastgpt/service/worker/tokenWorkerConfig';
describe('token worker config', () => {
beforeEach(() => {
osMock.availableParallelism.mockReset();
osMock.cpus.mockReset();
});
it('should cap token workers at 4 even when more CPU is available', () => {
osMock.availableParallelism.mockReturnValue(10);
expect(getTokenWorkerCount()).toBe(4);
});
it('should follow available CPU when fewer than 4 CPUs are available', () => {
osMock.availableParallelism.mockReturnValue(2);
expect(getTokenWorkerCount()).toBe(2);
});
it('should keep at least 1 token worker when CPU detection fails', () => {
osMock.availableParallelism.mockReturnValue(undefined);
osMock.cpus.mockReturnValue([]);
expect(getTokenWorkerCount()).toBe(1);
});
});
This source diff could not be displayed because it is too large. You can view the blob instead.
import {
type ChatCompletionContentPart,
type ChatCompletionCreateParams,
type ChatCompletionMessageParam,
type ChatCompletionTool
} from '@fastgpt/global/core/ai/llm/type';
import { ChatCompletionRequestMessageRoleEnum } from '@fastgpt/global/core/ai/constants';
import o200kTokenizer from 'gpt-tokenizer/encoding/o200k_base';
export type CountGptMessagesTokensParams = {
messages: ChatCompletionMessageParam[];
tools?: ChatCompletionTool[];
functionCall?: ChatCompletionCreateParams.Function[];
};
type TokenizerApi = {
countTokens: (
input: string,
options?: {
disallowedSpecial?: Set<string> | 'all';
allowedSpecial?: Set<string> | 'all';
}
) => number;
};
const tokenizer: TokenizerApi = o200kTokenizer;
const noDisallowedSpecial = { disallowedSpecial: new Set<string>() };
/**
* FastGPT 的 worker token 计数统一使用 GPT 现代模型的 o200k_base 编码。
* 该路径只做上下文预算和缺 usage 时的近似兜底;供应商返回 usage 时仍以 usage 为准。
*/
export const GPT_TOKENIZER_ENCODING = 'o200k_base';
type CountableContentPart = ChatCompletionContentPart | { type: 'refusal'; refusal: string };
/**
* 将多模态 content part 转成可计数文本。
*
* 这里不尝试复刻各家模型对图片、音频、文件的精确计费规则,只把会进入上下文或
* 明显影响输入规模的字段纳入估算;真实计费仍以模型供应商返回的 usage 为准。
*/
const contentPartToText = (part: CountableContentPart) => {
if (part.type === 'text') return part.text;
if (part.type === 'image_url') return part.image_url.url;
if (part.type === 'input_audio') return part.input_audio.data;
if (part.type === 'file')
return [part.file.filename, part.file.file_id, part.file.file_data].filter(Boolean).join(' ');
if (part.type === 'file_url') return [part.name, part.url].filter(Boolean).join(' ');
if (part.type === 'refusal') return part.refusal;
return '';
};
/**
* 统一把 OpenAI chat content 规整为字符串。
*
* 字符串 content 直接计数;数组 content 按 part 拼接,保持和旧方案一致的“近似预算”
* 语义,避免在不同消息类型间引入额外分隔符导致历史 token 预算明显漂移。
*/
const contentToText = (content: ChatCompletionMessageParam['content'] = '') => {
if (!content) return '';
if (typeof content === 'string') return content;
return (content as CountableContentPart[]).map(contentPartToText).join('');
};
const countTextTokens = (text: string) => {
try {
return tokenizer.countTokens(text, noDisallowedSpecial);
} catch {
// tokenizer 对极少数非法 special token 组合可能抛错,退回字符数保证计费链路不断。
return text.length;
}
};
/**
* 统计普通 prompt 文本 token 数。
*
* 该函数只在 token worker 内执行,用于知识库裁剪、embedding/rerank 兜底计费等
* 近似场景,统一按 o200k_base 估算。
*/
export const countPromptTokensInWorker = (
prompt: string | ChatCompletionContentPart[] | null | undefined = '',
role: '' | `${ChatCompletionRequestMessageRoleEnum}` = ''
) => {
const promptText =
typeof prompt === 'string' || !prompt ? prompt || '' : prompt.map(contentPartToText).join('');
const text = `${role}\n${promptText}`.trim();
// 兼容旧实现:只有传入 role 时才补 chat message 的固定结构开销。
const supplementaryToken = role ? 4 : 0;
return countTextTokens(text) + supplementaryToken;
};
const countToolsTokens = (tools?: ChatCompletionTool[] | ChatCompletionCreateParams.Function[]) => {
if (!tools || tools.length === 0) return 0;
// 旧方案也是把工具 schema 规整成紧凑文本后估算,避免格式化 JSON 的空白影响预算。
const toolText = JSON.stringify(tools)
.replace(/"/g, '')
.replace(/\n/g, '')
.replace(/( ){2,}/g, ' ');
return countTextTokens(toolText);
};
const getAssistantCallText = (message: ChatCompletionMessageParam) => {
if (message.role !== ChatCompletionRequestMessageRoleEnum.Assistant) return '';
// assistant 的 tool/function call 参数会进入模型上下文,需要和普通 content 一起计入。
const toolCallsText =
message.tool_calls
?.map((item) => `${item?.function?.name} ${item?.function?.arguments}`.trim())
?.join('') || '';
const functionCall = message.function_call;
const functionCallText = `${functionCall?.name || ''} ${functionCall?.arguments || ''}`.trim();
return `${toolCallsText}${functionCallText}`;
};
/**
* 在 token worker 内同步统计 Chat messages token 数。
*
* 这里保持旧实现的消息常数近似规则,只替换为更快的 GPT tokenizer。
* 主线程只通过 worker 调用该函数,避免主进程加载 tokenizer rank 常驻内存。
*/
export const countGptMessagesTokensInWorker = ({
messages,
tools,
functionCall
}: CountGptMessagesTokensParams) => {
return (
messages.reduce((sum, item, index) => {
// 只有最后一条消息的 reasoning_content 会继续影响后续上下文预算。
const reasoningText = index === messages.length - 1 ? item.reasoning_content || '' : '';
const contentPrompt = contentToText(item.content);
const callPrompt = getAssistantCallText(item);
const text = `${item.role}\n${reasoningText}${contentPrompt}${callPrompt}`.trim();
// 每条带 role 的 chat message 保留旧实现的固定结构开销,降低切换 tokenizer 的行为差异。
return sum + countTextTokens(text) + (item.role ? 4 : 0);
}, 0) +
countToolsTokens(tools) +
countToolsTokens(functionCall)
);
};
/* Only the token of gpt-3.5-turbo is used */
import { Tiktoken } from 'tiktoken/lite';
import cl100k_base from './cl100k_base.json';
import {
type ChatCompletionMessageParam,
type ChatCompletionContentPart,
type ChatCompletionCreateParams,
type ChatCompletionTool
} from '@fastgpt/global/core/ai/llm/type';
import { ChatCompletionRequestMessageRoleEnum } from '@fastgpt/global/core/ai/constants';
import { parentPort } from 'worker_threads';
import { getLogger, LogCategories } from '../../common/logger';
import { countGptMessagesTokensInWorker, countPromptTokensInWorker } from './count';
const enc = new Tiktoken(cl100k_base.bpe_ranks, cl100k_base.special_tokens, cl100k_base.pat_str);
const logger = getLogger(LogCategories.INFRA.WORKER);
/* count messages tokens */
type CountGptMessagesTokensWorkerPayload = {
id: string;
messages?: ChatCompletionMessageParam[];
messageGroups?: ChatCompletionMessageParam[][];
prompts?: (string | null | undefined)[];
tools?: ChatCompletionTool[];
functionCall?: ChatCompletionCreateParams.Function[];
};
/**
* Token 计数 worker 入口。
*
* 单条 messages、批量 messageGroups、批量 prompts 共用同一个 worker 文件,减少 worker
* 类型数量和初始化成本。批量请求在 worker 内同步 map,避免主线程为大量短文本反复
* postMessage,也保证返回顺序和输入顺序一致。
*/
parentPort?.on(
'message',
({
id,
messages,
messageGroups,
prompts,
tools,
functionCall
}: {
id: string;
messages: ChatCompletionMessageParam[];
tools?: ChatCompletionTool[];
functionCall?: ChatCompletionCreateParams.Function[];
}) => {
}: CountGptMessagesTokensWorkerPayload) => {
try {
/* count one prompt tokens */
const countPromptTokens = (
prompt: string | ChatCompletionContentPart[] | null | undefined = '',
role: '' | `${ChatCompletionRequestMessageRoleEnum}` = ''
) => {
const promptText = (() => {
if (!prompt) return '';
if (typeof prompt === 'string') return prompt;
let promptText = '';
prompt.forEach((item) => {
if (item.type === 'text') {
promptText += item.text;
} else if (item.type === 'image_url') {
promptText += item.image_url.url;
}
});
return promptText;
})();
const text = `${role}\n${promptText}`.trim();
try {
const encodeText = enc.encode(text);
const supplementaryToken = role ? 4 : 0;
return encodeText.length + supplementaryToken;
} catch (error) {
return text.length;
const data = (() => {
// 上下文裁剪会频繁计算多组 messages,批量放进一次 worker 消息能降低 IPC 开销。
if (messageGroups) {
return messageGroups.map((messages) => countGptMessagesTokensInWorker({ messages }));
}
};
const countToolsTokens = (
tools?: ChatCompletionTool[] | ChatCompletionCreateParams.Function[]
) => {
if (!tools || tools.length === 0) return 0;
const toolText = tools
? JSON.stringify(tools)
.replace('"', '')
.replace('\n', '')
.replace(/( ){2,}/g, ' ')
: '';
return enc.encode(toolText).length;
};
const total =
messages.reduce((sum, item, index) => {
// Evaluates the text of toolcall and functioncall
const functionCallPrompt = (() => {
let prompt = '';
if (item.role === ChatCompletionRequestMessageRoleEnum.Assistant) {
const toolCalls = item.tool_calls;
prompt +=
toolCalls
?.map((item) => `${item?.function?.name} ${item?.function?.arguments}`.trim())
?.join('') || '';
const functionCall = item.function_call;
prompt += `${functionCall?.name} ${functionCall?.arguments}`.trim();
}
return prompt;
})();
const contentPrompt = (() => {
if (!item.content) return '';
if (typeof item.content === 'string') return item.content;
return item.content
.map((item) => {
if (item.type === 'text') return item.text;
return '';
})
.join('');
})();
// Only the last message computed reasoning_content
const reasoningText = index === messages.length - 1 ? item.reasoning_content || '' : '';
// embedding/rerank 等路径多为纯文本 prompt,走轻量分支可少做 chat message 拼装。
if (prompts) {
return prompts.map((prompt) => countPromptTokensInWorker(prompt));
}
return (
sum +
countPromptTokens(`${reasoningText}${contentPrompt}${functionCallPrompt}`, item.role)
);
}, 0) +
countToolsTokens(tools) +
countToolsTokens(functionCall);
return countGptMessagesTokensInWorker({
messages: messages || [],
tools,
functionCall
});
})();
parentPort?.postMessage({
id,
type: 'success',
data: total
data
});
} catch (error) {
logger.error('Token count worker failed', { error });
parentPort?.postMessage({
id,
type: 'success',
data: 0
type: 'error',
data: error instanceof Error ? error.message : String(error)
});
}
}
......
import { getWorkerController, WorkerNameEnum } from './utils';
import { getLogger, LogCategories } from '../common/logger';
import { getTokenWorkerCount } from './tokenWorkerConfig';
const logger = getLogger(LogCategories.SYSTEM);
/**
* 服务启动时预热 token worker 池,避免首次聊天请求在业务路径上承担 tokenizer 初始化成本。
*
* 当前 worker 数量已经被限制为 min(cpu, 4),一次性并发创建的资源峰值可控;
* 因此不再分批预热,减少启动阶段的额外等待和中间日志噪音。
*/
export const preLoadWorker = async () => {
const start = Date.now();
const max = Math.min(Number(global.systemEnv?.tokenWorkers || 30), 500);
const max = getTokenWorkerCount();
const workerController = getWorkerController({
name: WorkerNameEnum.countGptMessagesTokens,
maxReservedThreads: max
});
// Batch size for concurrent loading
const batchSize = 5;
for (let i = 0; i < max; i += batchSize) {
const currentBatchSize = Math.min(batchSize, max - i);
const promises = [];
// Create batch of promises for concurrent loading
for (let j = 0; j < currentBatchSize; j++) {
const promise = (async () => {
const worker = workerController.createWorker();
await workerController.run({
workerId: worker.id,
messages: [
{
role: 'user',
content: '1'
}
]
});
})();
promises.push(promise);
}
// Wait for current batch to complete
await Promise.all(promises);
logger.debug('Worker preload batch completed', {
queuedWorkers: workerController.workerQueue.length,
batchSize: currentBatchSize,
max
});
}
await Promise.all(
Array.from({ length: max }, async () => {
const worker = workerController.createWorker();
// 用最小消息触发 worker 内 tokenizer 的首次加载,后续请求可复用已初始化 worker。
await workerController.run({
workerId: worker.id,
messages: [
{
role: 'user',
content: '1'
}
]
});
})
);
logger.info('Worker preload completed', {
queuedWorkers: workerController.workerQueue.length,
......
import { availableParallelism, cpus } from 'os';
/**
* Token 计算 worker 数量跟随当前运行环境可用 CPU 数,最多保留 4 个。
*
* tokenizer 会在每个 worker 内各自加载一份编码表,worker 过多会放大常驻内存;
* 这里固定为 min(cpu, 4),不再暴露配置项,避免部署环境误配过多 worker。
* availableParallelism 会优先考虑容器 CPU 配额,拿不到时再回退到物理 CPU 数。
*/
export const getTokenWorkerCount = () => {
const availableCpu = availableParallelism?.() || cpus().length || 1;
return Math.max(1, Math.min(availableCpu, 4));
};
......@@ -93,6 +93,9 @@ catalogs:
file-type:
specifier: 21.3.0
version: 21.3.0
gpt-tokenizer:
specifier: 3.4.0
version: 3.4.0
i18next:
specifier: 23.16.8
version: 23.16.8
......@@ -518,6 +521,9 @@ importers:
form-data:
specifier: ^4.0.4
version: 4.0.4
gpt-tokenizer:
specifier: 'catalog:'
version: 3.4.0
http-proxy-agent:
specifier: ^7.0.2
version: 7.0.2
......@@ -614,9 +620,6 @@ importers:
proxy-from-env:
specifier: ^1.1.0
version: 1.1.0
tiktoken:
specifier: 1.0.17
version: 1.0.17
turndown:
specifier: ^7.1.2
version: 7.2.0
......@@ -1462,9 +1465,6 @@ importers:
qs:
specifier: ^6.15.2
version: 6.15.2
tiktoken:
specifier: 1.0.17
version: 1.0.17
uuid:
specifier: ^9.0.1
version: 9.0.1
......@@ -9657,6 +9657,9 @@ packages:
resolution: {integrity: sha512-QLV1qeYSo5l13mQzWgP/y0LbMr5Plr5fJilgAIwgnwseproEbtNym8xpLsDzeZ6MWXgNE6kdWGBjdh3zT/Qerg==}
engines: {node: '>=20'}
gpt-tokenizer@3.4.0:
resolution: {integrity: sha512-wxFLnhIXTDjYebd9A9pGl3e31ZpSypbpIJSOswbgop5jLte/AsZVDvjlbEuVFlsqZixVKqbcoNmRlFDf6pz/UQ==}
graceful-fs@4.2.10:
resolution: {integrity: sha512-9ByhssR2fPVsNZj478qUUbKfmL0+t5BDVyjShtyZZLiK7ZDAArFFfopyOTj0M05wE2tJPisA4iTnnXl2YoPvOA==}
......@@ -13825,9 +13828,6 @@ packages:
through@2.3.8:
resolution: {integrity: sha512-w89qg7PI8wAdvX60bMDP+bFoD5Dvhm9oLheFp5O4a2QF0cSBGsBX4qZmadPMvVqlLJBBci+WqGGOAPvcDeNSVg==}
tiktoken@1.0.17:
resolution: {integrity: sha512-UuFHqpy/DxOfNiC3otsqbx3oS6jr5uKdQhB/CvDEroZQbVHt+qAK+4JbIooabUWKU9g6PpsFylNu9Wcg4MxSGA==}
timezones-list@3.1.0:
resolution: {integrity: sha512-PcDBt9tae330KTOIufK/wArTlJp+unuuRcG0EEu+4oLHZACHefKQyP2D51gMZID+urye92mHND60KRVuDDAmbA==}
......@@ -25062,6 +25062,8 @@ snapshots:
responselike: 4.0.2
type-fest: 4.41.0
gpt-tokenizer@3.4.0: {}
graceful-fs@4.2.10: {}
graceful-fs@4.2.11: {}
......@@ -30342,8 +30344,6 @@ snapshots:
through@2.3.8: {}
tiktoken@1.0.17: {}
timezones-list@3.1.0: {}
tiny-invariant@1.3.3: {}
......@@ -45,6 +45,7 @@ catalog:
dayjs: 1.11.19
express: ^4
file-type: 21.3.0
gpt-tokenizer: 3.4.0
i18next: 23.16.8
ipaddr.js: ^2.4.0
js-yaml: ^4.1.1
......
Subproject commit 128dc15d02ab52ce55a791c05cf6041cddd98271
Subproject commit b498edbc927f497f334aab9a3575de8e86f56001
......@@ -86,9 +86,6 @@ COPY --from=builder --chown=nextjs:nodejs /app/projects/app/.next/server/chunks
COPY --from=builder --chown=nextjs:nodejs /app/projects/app/worker /app/projects/app/worker
COPY --from=builder --chown=nextjs:nodejs /app/projects/app/worker /app/worker
# copy standload packages
COPY --from=maindeps /app/node_modules/tiktoken ./node_modules/tiktoken
RUN rm -rf ./node_modules/tiktoken/encoders
COPY --from=maindeps /app/node_modules/@zilliz/milvus2-sdk-node ./node_modules/@zilliz/milvus2-sdk-node
# copy package.json to version file
......
......@@ -9,7 +9,6 @@
"vectorMaxProcess": 10, // 向量处理线程数量
"qaMaxProcess": 10, // 问答拆分线程数量
"vlmMaxProcess": 10, // 图片理解模型最大处理进程
"tokenWorkers": 30, // Token 计算线程保持数,会持续占用内存,不能设置太大。
"hnswEfSearch": 100, // 向量搜索参数,仅对 PG OB 生效。越大,搜索越精确,但是速度越慢。设置为100,有99%+精度。
"hnswMaxScanTuples": 100000, // 向量搜索最大扫描数据量,仅对 PG生效。
"customPdfParse": {
......
......@@ -76,7 +76,6 @@ const nextConfig: NextConfig = {
'@node-rs/jieba',
'bullmq',
'@zilliz/milvus2-sdk-node',
'tiktoken',
'@opentelemetry/api-logs',
'@mariozechner/pi-agent-core',
'@mariozechner/pi-ai'
......
import type { ApiRequestProps, ApiResponseType } from '@fastgpt/service/type/next';
import { NextAPI } from '@/service/middleware/entry';
import { authCert } from '@fastgpt/service/support/permission/auth/common';
import { type ChatCompletionMessageParam } from '@fastgpt/global/core/ai/llm/type';
import { countGptMessagesTokens } from '@fastgpt/service/common/string/tiktoken';
export type tokenQuery = {};
export type tokenBody = {
messages: ChatCompletionMessageParam[];
};
export type tokenResponse = {};
async function handler(
req: ApiRequestProps<tokenBody, tokenQuery>,
res: ApiResponseType<any>
): Promise<tokenResponse> {
const start = Date.now();
await authCert({ req, authRoot: true });
const tokens = await countGptMessagesTokens(req.body.messages);
return {
tokens,
time: Date.now() - start,
memory: process.memoryUsage()
};
}
export default NextAPI(handler);
export const config = {
api: {
bodyParser: {
sizeLimit: '20mb'
},
responseLimit: '20mb'
}
};
......@@ -32,7 +32,6 @@
"lodash": "catalog:",
"moment": "^2.30.1",
"qs": "^6.15.2",
"tiktoken": "1.0.17",
"uuid": "^9.0.1",
"zod": "catalog:"
},
......
......@@ -76,7 +76,6 @@ const nextConfig: NextConfig = {
'pg',
'bullmq',
'@zilliz/milvus2-sdk-node',
'tiktoken',
'@opentelemetry/api-logs'
],
experimental: {
......
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