Commit 05bb1979 by Archer Committed by GitHub

V4.14.9 features (#6599)

* fix: image read and json error (Agent) (#6502)

* fix:
1.image read
2.JSON parsing error

* dataset cite and pause

* perf: plancall second parse

* add test

---------

Co-authored-by: archer <545436317@qq.com>

* master message

* remove invalid code

* feat(sre): integrate traces, logs, metrics into one sdk (#6580)

* fix: image read and json error (Agent) (#6502)

* fix:
1.image read
2.JSON parsing error

* dataset cite and pause

* perf: plancall second parse

* add test

---------

Co-authored-by: archer <545436317@qq.com>

* master message

* wip: otel sdk

* feat(sre): integrate traces, logs, metrics into one sdk

* fix(sre): use SpanStatusCode constants

* fix(sre): clarify step memory measurement

* update package

* fix: ts

---------

Co-authored-by: YeYuheng <57035043+YYH211@users.noreply.github.com>
Co-authored-by: archer <545436317@qq.com>

* doc

* sandbox in agent (#6579)

* doc

* update template

* fix: pr

* fix: sdk package

* update lock

* update next

* update dockerfile

* dockerfile

* dockerfile

* update sdk version

* update dockerefile

* version

---------

Co-authored-by: YeYuheng <57035043+YYH211@users.noreply.github.com>
Co-authored-by: Ryo <whoeverimf5@gmail.com>
parent 7a660139
......@@ -78,6 +78,14 @@ mock 对应的 API 请求进行测试。
2. 如果测试不通过,则根据错误信息检查代码逻辑或者测试用例。
3. 如需二次修改,则回到”二、测例编写“。
## 单测包含哪些场景
1. 基础场景
2. 复杂场景
3. 边界值
4. 安全边界情况(死循环、系统崩溃、超大数据等)
5. 异常场景
## 常用命令
```shell
......
......@@ -26,9 +26,10 @@ AGENT_SANDBOX_SEALOS_TOKEN=
## 🚀 新增内容
1. 新增 AI 虚拟机功能,可以给 AI 挂载一个虚拟机工具进行更丰富的操作。
2. 封装 logger sdk。
3. 更新知识库单个数据时,同步更新 collection 更新时间。
4. 表单输入文件时,支持打开文件进行预览。
2. AgentV2 上下文适配暂停态。
3. 封装 logger sdk。增加 Metrics 追踪。
4. 更新知识库单个数据时,同步更新 collection 更新时间。
5. 表单输入文件时,支持打开文件进行预览。
## ⚙️ 优化
......@@ -54,3 +55,4 @@ AGENT_SANDBOX_SEALOS_TOKEN=
13. 系统工具集不显示版本
14. 修复视频音频自定义文件类型流程开始无文件链接变量
15. 用户输入框消息不转义成 markdown 格式
16. 修复 AgentV2 部分上下文错误。
......@@ -77,6 +77,7 @@ description: FastGPT Toc
- [/en/docs/openapi/share](/en/docs/openapi/share)
- [/en/docs/self-host/config/json](/en/docs/self-host/config/json)
- [/en/docs/self-host/config/model/intro](/en/docs/self-host/config/model/intro)
- [/en/docs/self-host/config/model/minimax](/en/docs/self-host/config/model/minimax)
- [/en/docs/self-host/config/model/siliconCloud](/en/docs/self-host/config/model/siliconCloud)
- [/en/docs/self-host/config/object-storage](/en/docs/self-host/config/object-storage)
- [/en/docs/self-host/config/signoz](/en/docs/self-host/config/signoz)
......
......@@ -77,6 +77,7 @@ description: FastGPT 文档目录
- [/docs/openapi/share](/docs/openapi/share)
- [/docs/self-host/config/json](/docs/self-host/config/json)
- [/docs/self-host/config/model/intro](/docs/self-host/config/model/intro)
- [/docs/self-host/config/model/minimax](/docs/self-host/config/model/minimax)
- [/docs/self-host/config/model/siliconCloud](/docs/self-host/config/model/siliconCloud)
- [/docs/self-host/config/object-storage](/docs/self-host/config/object-storage)
- [/docs/self-host/config/signoz](/docs/self-host/config/signoz)
......
......@@ -149,6 +149,8 @@
"document/content/docs/self-host/config/json.mdx": "2026-03-03T17:39:47+08:00",
"document/content/docs/self-host/config/model/intro.en.mdx": "2026-03-19T14:09:03+08:00",
"document/content/docs/self-host/config/model/intro.mdx": "2026-03-19T14:09:03+08:00",
"document/content/docs/self-host/config/model/minimax.en.mdx": "2026-03-19T09:32:57-05:00",
"document/content/docs/self-host/config/model/minimax.mdx": "2026-03-19T09:32:57-05:00",
"document/content/docs/self-host/config/model/siliconCloud.en.mdx": "2026-03-19T14:09:03+08:00",
"document/content/docs/self-host/config/model/siliconCloud.mdx": "2026-03-19T14:09:03+08:00",
"document/content/docs/self-host/config/object-storage.en.mdx": "2026-03-03T17:39:47+08:00",
......@@ -236,7 +238,7 @@
"document/content/docs/self-host/upgrading/4-14/4148.mdx": "2026-03-09T17:39:53+08:00",
"document/content/docs/self-host/upgrading/4-14/41481.en.mdx": "2026-03-09T12:02:02+08:00",
"document/content/docs/self-host/upgrading/4-14/41481.mdx": "2026-03-09T17:39:53+08:00",
"document/content/docs/self-host/upgrading/4-14/4149.mdx": "2026-03-19T14:09:03+08:00",
"document/content/docs/self-host/upgrading/4-14/4149.mdx": "2026-03-20T22:01:38+08:00",
"document/content/docs/self-host/upgrading/outdated/40.en.mdx": "2026-03-03T17:39:47+08:00",
"document/content/docs/self-host/upgrading/outdated/40.mdx": "2026-03-03T17:39:47+08:00",
"document/content/docs/self-host/upgrading/outdated/41.en.mdx": "2026-03-03T17:39:47+08:00",
......@@ -377,8 +379,8 @@
"document/content/docs/self-host/upgrading/outdated/499.mdx": "2026-03-03T17:39:47+08:00",
"document/content/docs/self-host/upgrading/upgrade-intruction.en.mdx": "2026-03-03T17:39:47+08:00",
"document/content/docs/self-host/upgrading/upgrade-intruction.mdx": "2026-03-03T17:39:47+08:00",
"document/content/docs/toc.en.mdx": "2026-03-19T14:09:03+08:00",
"document/content/docs/toc.mdx": "2026-03-19T14:09:03+08:00",
"document/content/docs/toc.en.mdx": "2026-03-20T21:57:22+08:00",
"document/content/docs/toc.mdx": "2026-03-20T21:57:22+08:00",
"document/content/docs/use-cases/app-cases/dalle3.en.mdx": "2026-02-26T22:14:30+08:00",
"document/content/docs/use-cases/app-cases/dalle3.mdx": "2025-07-23T21:35:03+08:00",
"document/content/docs/use-cases/app-cases/english_essay_correction_bot.en.mdx": "2026-02-26T22:14:30+08:00",
......
......@@ -7,6 +7,7 @@ export const topAgentParamsSchema = z.object({
systemPrompt: z.string().nullish(),
selectedTools: z.array(z.string()).nullish(),
selectedDatasets: z.array(z.string()).nullish(),
fileUpload: z.boolean().nullish()
fileUpload: z.boolean().nullish(),
enableSandbox: z.boolean().nullish()
});
export type TopAgentParamsType = z.infer<typeof topAgentParamsSchema>;
import { configureLoggerFromEnv, disposeLogger, getLogger } from '@fastgpt-sdk/logger';
import { configureLoggerFromEnv, disposeLogger, getLogger } from '@fastgpt-sdk/otel/logger';
import { env } from '../../env';
export async function configureLogger() {
......
export { configureLogger, disposeLogger, getLogger } from './client';
export { withContext, withCategoryPrefix } from '@fastgpt-sdk/logger';
export { withContext, withCategoryPrefix } from '@fastgpt-sdk/otel/logger';
export { LogCategories } from './categories';
export type { LogCategory } from './categories';
import { configureMetricsFromEnv, disposeMetrics, getMeter } from '@fastgpt-sdk/otel/metrics';
import { env } from '../../env';
export async function configureMetrics() {
await configureMetricsFromEnv({
env,
defaultServiceName: 'fastgpt-client',
defaultMeterName: 'fastgpt-client'
});
}
export { disposeMetrics, getMeter };
export { configureMetrics, disposeMetrics, getMeter } from './client';
import { jsonRes } from '../response';
import type { NextApiRequest, NextApiResponse } from 'next';
import { SpanStatusCode } from '@opentelemetry/api';
import { withNextCors } from './cors';
import { type ApiRequestProps } from '../../type/next';
import { getLogger, LogCategories, withContext } from '../logger';
import { setSpanError, withActiveSpan } from '../tracing';
import { ZodError } from 'zod';
import { randomUUID } from 'crypto';
......@@ -24,7 +26,6 @@ export const NextEntry = ({
const requestLogger = getLogger(LogCategories.HTTP.REQUEST);
const responseLogger = getLogger(LogCategories.HTTP.RESPONSE);
const errorLogger = getLogger(LogCategories.HTTP.ERROR);
const url = req.url || '';
const method = req.method?.toUpperCase() || '';
......@@ -32,75 +33,105 @@ export const NextEntry = ({
const userAgent = req.headers['user-agent'];
const contentLength = req.headers['content-length'];
return withContext({ requestId }, async () => {
requestLogger.info(`[${method}] ${url}`, {
verbose: false,
requestId,
method,
url,
ip,
userAgent,
contentLength
});
let responseLogged = false;
const logResponse = (event: 'request-finish' | 'request-close') => {
if (responseLogged) return;
responseLogged = true;
const durationMs = Date.now() - start;
const httpStatusCode = res.statusCode;
responseLogger.info(`[${method}] ${url} - ${httpStatusCode} in ${durationMs}ms`, {
verbose: false,
requestId,
method,
httpStatusCode,
event
});
};
res.once('finish', () => logResponse('request-finish'));
res.once('close', () => logResponse('request-close'));
try {
await Promise.all([
withNextCors(req, res),
...beforeCallback.map((item) => item(req, res))
]);
let response = null;
for await (const handler of args) {
response = await handler(req, res);
if (res.writableFinished) {
break;
return withContext({ requestId }, async () =>
withActiveSpan(
{
name: `http.request ${method || 'UNKNOWN'} ${url || '/'}`,
tracerName: 'fastgpt.http',
attributes: {
'fastgpt.request.id': requestId,
'http.request.method': method,
'url.full': url,
'client.address': Array.isArray(ip) ? ip.join(',') : ip,
'user_agent.original': userAgent,
'http.request.body.size': contentLength
}
}
const contentType = res.getHeader('Content-Type');
if ((!contentType || contentType === 'application/json') && !res.writableFinished) {
return jsonRes(res, {
code: 200,
data: response
},
async (span) => {
requestLogger.info(`[${method}] ${url}`, {
verbose: false,
requestId,
method,
url,
ip,
userAgent,
contentLength
});
}
} catch (error) {
// Handle Zod validation errors
if (error instanceof ZodError) {
return jsonRes(res, {
code: 400,
message: 'Data validation error',
error,
url: req.url
});
}
return jsonRes(res, {
code: 500,
error,
url: req.url
});
}
});
let responseLogged = false;
const logResponse = (event: 'request-finish' | 'request-close') => {
if (responseLogged) return;
responseLogged = true;
const durationMs = Date.now() - start;
const httpStatusCode = res.statusCode;
responseLogger.info(`[${method}] ${url} - ${httpStatusCode} in ${durationMs}ms`, {
verbose: false,
requestId,
method,
httpStatusCode,
event
});
};
res.once('finish', () => logResponse('request-finish'));
res.once('close', () => logResponse('request-close'));
try {
await Promise.all([
withNextCors(req, res),
...beforeCallback.map((item) => item(req, res))
]);
let response = null;
for await (const handler of args) {
response = await handler(req, res);
if (res.writableFinished) {
break;
}
}
const contentType = res.getHeader('Content-Type');
if ((!contentType || contentType === 'application/json') && !res.writableFinished) {
const jsonResponse = await jsonRes(res, {
code: 200,
data: response
});
span.setAttribute('http.response.status_code', res.statusCode);
return jsonResponse;
}
span.setAttribute('http.response.status_code', res.statusCode);
} catch (error) {
// Handle Zod validation errors
if (error instanceof ZodError) {
span.setAttribute('http.response.status_code', 400);
span.setStatus({
code: SpanStatusCode.ERROR,
message: 'Data validation error'
});
return jsonRes(res, {
code: 400,
message: 'Data validation error',
error,
url: req.url
});
}
span.setAttribute('http.response.status_code', 500);
setSpanError(span, error);
return jsonRes(res, {
code: 500,
error,
url: req.url
});
}
}
)
);
};
};
};
import { getErrText } from '@fastgpt/global/common/error/utils';
import { SpanStatusCode } from '@opentelemetry/api';
import {
configureTracingFromEnv,
disposeTracing,
getCurrentSpanContext,
getTracer
} from '@fastgpt-sdk/otel/tracing';
import { withContext } from '../logger';
import { env } from '../../env';
type SpanAttributeValue = string | number | boolean;
type SpanStatusLike = {
code?: number;
message?: string;
};
type TracerLike = ReturnType<typeof getTracer>;
type SpanLike = ReturnType<TracerLike['startSpan']>;
export type TraceLogContext = {
traceId: string;
spanId: string;
};
export type ActiveSpanOptions = {
name: string;
tracer?: TracerLike;
tracerName?: string;
attributes?: Record<string, unknown>;
};
function normalizeAttributes(attributes?: Record<string, unknown>) {
if (!attributes) return;
const normalized: Record<string, SpanAttributeValue> = {};
Object.entries(attributes).forEach(([key, value]) => {
if (value === undefined || value === null) return;
if (typeof value === 'string' || typeof value === 'number' || typeof value === 'boolean') {
normalized[key] = value satisfies SpanAttributeValue;
return;
}
});
return Object.keys(normalized).length > 0 ? normalized : undefined;
}
export async function configureTracing() {
await configureTracingFromEnv({
env,
defaultServiceName: 'fastgpt-client',
defaultTracerName: 'fastgpt-client',
defaultSampleRatio: env.TRACING_OTEL_SAMPLE_RATIO
});
}
export function getTraceLogContext(): TraceLogContext | undefined {
const spanContext = getCurrentSpanContext();
if (!spanContext) return;
return {
traceId: spanContext.traceId,
spanId: spanContext.spanId
};
}
export function setSpanError(
span: SpanLike,
error: unknown,
extraStatus?: Partial<SpanStatusLike>
) {
span.recordException(error instanceof Error ? error : new Error(getErrText(error)));
span.setStatus({
code: SpanStatusCode.ERROR,
message: extraStatus?.message ?? getErrText(error)
});
}
export async function withActiveSpan<T>(
options: ActiveSpanOptions,
callback: (span: SpanLike) => Promise<T> | T
): Promise<T> {
const tracer = options.tracer ?? getTracer(options.tracerName);
return tracer.startActiveSpan(
options.name,
{
attributes: normalizeAttributes(options.attributes)
},
async (span: SpanLike) => {
const spanContext = span.spanContext();
return withContext(
{
traceId: spanContext.traceId,
spanId: spanContext.spanId
},
async () => {
try {
return await callback(span);
} catch (error) {
setSpanError(span, error);
throw error;
} finally {
span.end();
}
}
);
}
);
}
export { disposeTracing, getCurrentSpanContext, getTracer };
export {
configureTracing,
disposeTracing,
getCurrentSpanContext,
getTraceLogContext,
getTracer,
setSpanError,
withActiveSpan
} from './client';
export type { ActiveSpanOptions, TraceLogContext } from './client';
......@@ -160,6 +160,7 @@ export const dispatchTopAgent = async (
tools, // 从 execution_plan 提取
datasets: filterDatasets,
fileUploadEnabled: responseJson.resources?.system_features?.file_upload?.enabled || false,
enableSandboxEnabled: responseJson.resources?.system_features?.sandbox?.enabled || false,
executionPlan: responseJson.execution_plan // 保存原始 execution_plan
});
......
......@@ -34,6 +34,12 @@ export const getPrompt = ({
);
}
if (metadata.enableSandbox !== undefined && metadata.enableSandbox !== null) {
sections.push(
`**虚拟机**: ${metadata.enableSandbox ? '搭建者已启用虚拟机功能' : '搭建者已禁用虚拟机功能'}`
);
}
if (sections.length === 0) return '';
return `
......@@ -337,6 +343,8 @@ ${resourceList}
- 系统功能判断:
* 是否需要用户的私有文件?→ 启用 file_upload
* 数据能否通过工具获取?→ 不需要 file_upload
* 是否需要执行代码或数据处理(如运行 Python 脚本、复杂计算、数据转换)?→ 启用 sandbox
* 任务仅需 LLM 推理和工具调用,无需执行任意代码?→ 不需要 sandbox
🔧 第四层:资源整合
- 收集所有需要的工具、知识库和系统功能
......@@ -377,6 +385,10 @@ ${resourceList}
"file_upload": {
"enabled": true/false,
"purpose": "说明原因(enabled=true时必填)"
},
"sandbox": {
"enabled": true/false,
"purpose": "说明为何需要虚拟机执行能力(enabled=true时必填,适用于代码执行、数据处理等场景)"
}
}
}
......@@ -395,6 +407,8 @@ ${resourceList}
- resources: 资源配置对象,仅包含系统功能配置
* system_features.file_upload.enabled: 是否需要文件上传(必填)
* system_features.file_upload.purpose: 为什么需要(enabled=true时必填)
* system_features.sandbox.enabled: 是否需要虚拟机执行能力(可选,适用于代码执行、数据处理场景)
* system_features.sandbox.purpose: 为什么需要虚拟机(enabled=true时必填)
<execution_plan_design>
**执行计划设计**:
......
......@@ -27,6 +27,7 @@ export const TopAgentFormDataSchema = z.object({
tools: z.array(z.string()).optional().default([]),
datasets: z.array(SelectedDatasetSchema).optional().default([]),
fileUploadEnabled: z.boolean().optional().default(false),
enableSandboxEnabled: z.boolean().optional().default(false),
executionPlan: z.any().optional()
});
export type TopAgentFormDataType = z.infer<typeof TopAgentFormDataSchema>;
......@@ -47,10 +48,19 @@ export const TopAgentGenerationAnswerSchema = z.object({
execution_plan: ExecutionPlanSchema.optional(),
resources: z.object({
system_features: z.object({
file_upload: z.object({
enabled: z.boolean(),
purpose: z.string().optional()
})
file_upload: z
.object({
enabled: z.boolean(),
purpose: z.string().optional()
})
.optional()
.default({ enabled: false }),
sandbox: z
.object({
enabled: z.boolean()
})
.optional()
.default({ enabled: false })
})
})
});
......
......@@ -65,7 +65,8 @@ ${tool}
${dataset}
### 系统功能
- **file_upload**: 文件上传功能 (enabled, purpose, file_types)
- **file_upload**: 文件上传功能,允许用户在对话中上传文件,让 Agent 读取私有文件内容
- **sandbox**: 虚拟机执行环境,为 Agent 提供代码运行能力(Python、Shell 等),适用于数据处理、科学计算、代码执行等场景
`;
};
......@@ -109,10 +110,12 @@ ${dataset}
})
]);
const allTools = [...systemTools, ...myTools];
const fileReadInfo = systemSubInfo[SubAppIds.fileRead];
const fileReadTool = `- **${SubAppIds.fileRead}** [工具]: ${parseI18nString(fileReadInfo.name, lang)} - ${fileReadInfo.toolDescription}`;
allTools.push(fileReadTool);
const builtinTools = [SubAppIds.fileRead, SubAppIds.sandboxTool].map((id) => {
const info = systemSubInfo[id];
return `- **${id}** [工具]: ${parseI18nString(info.name, lang)} - ${info.toolDescription}`;
});
const allTools = [...systemTools, ...myTools, ...builtinTools];
return {
resourceList: getPrompt({
......
......@@ -89,6 +89,7 @@ export const dispatchRunAgent = async (props: DispatchAgentModuleProps): Promise
userChatInput, // 本次任务的输入
history = 6,
fileUrlList: fileLinks,
aiChatVision = true,
agent_selectedTools: selectedTools = [],
// Dataset search configuration
agent_datasetParams: datasetParams,
......
import { getMeter } from '../../common/metrics';
type MetricAttributeValue = string | number | boolean;
type MetricAttributes = Record<string, MetricAttributeValue>;
export type WorkflowStepMetricAttributes = {
workflowId?: string;
workflowName?: string;
nodeId: string;
nodeName?: string;
nodeType: string;
mode?: string;
};
type ProcessSnapshot = {
rss: number;
heapUsed: number;
external: number;
arrayBuffers: number;
cpuUser: number;
cpuSystem: number;
};
type StepObservationState = {
startedAt: bigint;
startSnapshot: ProcessSnapshot;
hadOverlapAtStart: boolean;
overlapVersionAtStart: number;
};
function normalizeAttributes(attributes: Record<string, unknown>): MetricAttributes {
const normalized: MetricAttributes = {};
Object.entries(attributes).forEach(([key, value]) => {
if (value === undefined || value === null) return;
if (typeof value === 'string' || typeof value === 'number' || typeof value === 'boolean') {
normalized[key] = value;
}
});
return normalized;
}
function toMetricAttributes(
attributes: WorkflowStepMetricAttributes,
extras?: Record<string, unknown>
) {
return normalizeAttributes({
workflow_id: attributes.workflowId,
workflow_name: attributes.workflowName,
node_id: attributes.nodeId,
node_name: attributes.nodeName,
node_type: attributes.nodeType,
mode: attributes.mode,
...extras
});
}
function takeProcessSnapshot(): ProcessSnapshot {
const memory = process.memoryUsage();
const cpu = process.cpuUsage();
return {
rss: memory.rss,
heapUsed: memory.heapUsed,
external: memory.external,
arrayBuffers: memory.arrayBuffers,
cpuUser: cpu.user,
cpuSystem: cpu.system
};
}
let activeWorkflowStepCount = 0;
let overlapVersion = 0;
function beginStepObservation(): StepObservationState {
const state: StepObservationState = {
startedAt: process.hrtime.bigint(),
startSnapshot: takeProcessSnapshot(),
hadOverlapAtStart: activeWorkflowStepCount > 0,
overlapVersionAtStart: overlapVersion
};
activeWorkflowStepCount += 1;
if (activeWorkflowStepCount > 1) {
overlapVersion += 1;
}
return state;
}
const meter = getMeter('fastgpt.workflow');
const prefix = 'fastgpt.workflow';
const stepDuration = meter.createHistogram(`${prefix}.step.duration`, {
description: 'Workflow step execution duration',
unit: 'ms'
});
const stepExecutions = meter.createCounter(`${prefix}.step.executions`, {
description: 'Workflow step execution count'
});
const stepActive = meter.createUpDownCounter(`${prefix}.step.active`, {
description: 'Workflow steps currently executing'
});
const stepCpuUserTime = meter.createHistogram(`${prefix}.step.cpu.user_time`, {
description: 'Workflow step user CPU time',
unit: 'us'
});
const stepCpuSystemTime = meter.createHistogram(`${prefix}.step.cpu.system_time`, {
description: 'Workflow step system CPU time',
unit: 'us'
});
const stepMemoryRssStart = meter.createHistogram(`${prefix}.step.memory.rss_start`, {
description: 'Workflow process RSS memory snapshot at step start',
unit: 'By'
});
const stepMemoryHeapUsedStart = meter.createHistogram(`${prefix}.step.memory.heap_used_start`, {
description: 'Workflow process heap used memory snapshot at step start',
unit: 'By'
});
const stepMemoryExternalStart = meter.createHistogram(`${prefix}.step.memory.external_start`, {
description: 'Workflow process external memory snapshot at step start',
unit: 'By'
});
const stepMemoryArrayBuffersStart = meter.createHistogram(
`${prefix}.step.memory.array_buffers_start`,
{
description: 'Workflow process array buffer memory snapshot at step start',
unit: 'By'
}
);
const stepMemoryRss = meter.createHistogram(`${prefix}.step.memory.rss`, {
description: 'Workflow process RSS memory snapshot at step end',
unit: 'By'
});
const stepMemoryHeapUsed = meter.createHistogram(`${prefix}.step.memory.heap_used`, {
description: 'Workflow process heap used memory snapshot at step end',
unit: 'By'
});
const stepMemoryExternal = meter.createHistogram(`${prefix}.step.memory.external`, {
description: 'Workflow process external memory snapshot at step end',
unit: 'By'
});
const stepMemoryArrayBuffers = meter.createHistogram(`${prefix}.step.memory.array_buffers`, {
description: 'Workflow process array buffer memory snapshot at step end',
unit: 'By'
});
const stepMemoryRssGrowth = meter.createHistogram(`${prefix}.step.memory.rss_growth`, {
description: 'Workflow process RSS memory growth during non-overlapping step execution',
unit: 'By'
});
const stepMemoryHeapUsedGrowth = meter.createHistogram(`${prefix}.step.memory.heap_used_growth`, {
description: 'Workflow process heap used memory growth during non-overlapping step execution',
unit: 'By'
});
const stepMemoryExternalGrowth = meter.createHistogram(`${prefix}.step.memory.external_growth`, {
description: 'Workflow process external memory growth during non-overlapping step execution',
unit: 'By'
});
export async function observeWorkflowStep<T>(
attributes: WorkflowStepMetricAttributes,
fn: () => Promise<T> | T
): Promise<T> {
const observationState = beginStepObservation();
const baseAttributes = toMetricAttributes(attributes);
stepActive.add(1, baseAttributes);
try {
const result = await fn();
recordWorkflowStepEnd(attributes, observationState, 'ok', baseAttributes);
return result;
} catch (error) {
recordWorkflowStepEnd(attributes, observationState, 'error', baseAttributes);
throw error;
}
}
function recordWorkflowStepEnd(
attributes: WorkflowStepMetricAttributes,
observationState: StepObservationState,
status: 'ok' | 'error',
baseAttributes: MetricAttributes
) {
const endSnapshot = takeProcessSnapshot();
const metricAttributes = toMetricAttributes(attributes, { status });
const stepOverlap =
observationState.hadOverlapAtStart || observationState.overlapVersionAtStart !== overlapVersion;
const memoryAttributes = toMetricAttributes(attributes, {
status,
memory_scope: 'process',
memory_attribution: stepOverlap ? 'best_effort' : 'exclusive',
step_overlap: stepOverlap
});
const durationMs = Number(process.hrtime.bigint() - observationState.startedAt) / 1_000_000;
stepDuration.record(durationMs, metricAttributes);
stepExecutions.add(1, metricAttributes);
stepCpuUserTime.record(
Math.max(0, endSnapshot.cpuUser - observationState.startSnapshot.cpuUser),
metricAttributes
);
stepCpuSystemTime.record(
Math.max(0, endSnapshot.cpuSystem - observationState.startSnapshot.cpuSystem),
metricAttributes
);
stepMemoryRssStart.record(observationState.startSnapshot.rss, memoryAttributes);
stepMemoryHeapUsedStart.record(observationState.startSnapshot.heapUsed, memoryAttributes);
stepMemoryExternalStart.record(observationState.startSnapshot.external, memoryAttributes);
stepMemoryArrayBuffersStart.record(observationState.startSnapshot.arrayBuffers, memoryAttributes);
stepMemoryRss.record(endSnapshot.rss, memoryAttributes);
stepMemoryHeapUsed.record(endSnapshot.heapUsed, memoryAttributes);
stepMemoryExternal.record(endSnapshot.external, memoryAttributes);
stepMemoryArrayBuffers.record(endSnapshot.arrayBuffers, memoryAttributes);
if (!stepOverlap && endSnapshot.rss > observationState.startSnapshot.rss) {
stepMemoryRssGrowth.record(
endSnapshot.rss - observationState.startSnapshot.rss,
memoryAttributes
);
}
if (!stepOverlap && endSnapshot.heapUsed > observationState.startSnapshot.heapUsed) {
stepMemoryHeapUsedGrowth.record(
endSnapshot.heapUsed - observationState.startSnapshot.heapUsed,
memoryAttributes
);
}
if (!stepOverlap && endSnapshot.external > observationState.startSnapshot.external) {
stepMemoryExternalGrowth.record(
endSnapshot.external - observationState.startSnapshot.external,
memoryAttributes
);
}
activeWorkflowStepCount = Math.max(0, activeWorkflowStepCount - 1);
stepActive.add(-1, baseAttributes);
}
......@@ -22,7 +22,17 @@ export const env = createEnv({
LOG_ENABLE_OTEL: BoolSchema.default(false),
LOG_OTEL_LEVEL: LogLevelSchema.default('info'),
LOG_OTEL_SERVICE_NAME: z.string().default('fastgpt-client'),
LOG_OTEL_URL: z.string().url().optional()
LOG_OTEL_URL: z.url().optional(),
METRICS_ENABLE_OTEL: BoolSchema.default(false),
METRICS_EXPORT_INTERVAL: z.coerce.number().int().positive().default(15000),
METRICS_OTEL_SERVICE_NAME: z.string().default('fastgpt-client'),
METRICS_OTEL_URL: z.url().optional(),
TRACING_ENABLE_OTEL: BoolSchema.default(false),
TRACING_OTEL_SERVICE_NAME: z.string().default('fastgpt-client'),
TRACING_OTEL_URL: z.url().optional(),
TRACING_OTEL_SAMPLE_RATIO: z.coerce.number().min(0).max(1).optional()
},
emptyStringAsUndefined: true,
runtimeEnv: process.env,
......
......@@ -8,13 +8,14 @@
},
"dependencies": {
"@apidevtools/json-schema-ref-parser": "^11.7.2",
"@fastgpt-sdk/sandbox-adapter": "^0.0.22",
"@fastgpt-sdk/sandbox-adapter": "^0.0.27",
"@fastgpt-sdk/otel": "catalog:",
"@fastgpt-sdk/storage": "catalog:",
"@fastgpt-sdk/logger": "catalog:",
"@fastgpt/global": "workspace:*",
"@maxmind/geoip2-node": "^6.3.4",
"@modelcontextprotocol/sdk": "catalog:",
"@node-rs/jieba": "2.0.1",
"@opentelemetry/api": "^1.9.0",
"@t3-oss/env-core": "0.13.10",
"@xmldom/xmldom": "^0.8.10",
"@zilliz/milvus2-sdk-node": "2.4.10",
......@@ -36,8 +37,8 @@
"ioredis": "^5.6.0",
"joplin-turndown-plugin-gfm": "^1.0.12",
"json5": "catalog:",
"jsonrepair": "^3.0.0",
"jsonpath-plus": "^10.3.0",
"jsonrepair": "^3.0.0",
"jsonwebtoken": "^9.0.2",
"lodash": "catalog:",
"mammoth": "^1.11.0",
......
......@@ -7,7 +7,7 @@
*/
import type { CSSProperties } from 'react';
import { useEffect, useMemo, useState, useTransition } from 'react';
import { useEffect, useMemo, useState, useTransition, useRef } from 'react';
import { LexicalComposer } from '@lexical/react/LexicalComposer';
import { PlainTextPlugin } from '@lexical/react/LexicalPlainTextPlugin';
import { RichTextPlugin } from '@lexical/react/LexicalRichTextPlugin';
......@@ -33,7 +33,7 @@ import type { FormPropsType } from './type';
import { type EditorVariableLabelPickerType, type EditorVariablePickerType } from './type';
import { getNanoid } from '@fastgpt/global/common/string/tools';
import FocusPlugin from './plugins/FocusPlugin';
import { textToEditorState } from './utils';
import { textToEditorState, editorStateToText } from './utils';
import { MaxLengthPlugin } from './plugins/MaxLengthPlugin';
import { VariableLabelNode } from './plugins/VariableLabelPlugin/node';
import VariableLabelPlugin from './plugins/VariableLabelPlugin';
......@@ -145,6 +145,7 @@ export default function Editor({
const [_, startSts] = useTransition();
const [focus, setFocus] = useState(false);
const [scrollHeight, setScrollHeight] = useState(0);
const editorOutputRef = useRef(value);
const initialConfig = {
namespace: isRichText ? 'richPromptEditor' : 'promptEditor',
......@@ -164,7 +165,7 @@ export default function Editor({
};
useDeepCompareEffect(() => {
if (focus) return;
if (focus && value === editorOutputRef.current) return;
setKey(getNanoid(6));
}, [value, variables, variableLabels, skillOption, selectedSkills]);
......@@ -256,6 +257,7 @@ export default function Editor({
<OnBlurPlugin onBlur={onBlur} />
<OnChangePlugin
onChange={(editorState, editor) => {
editorOutputRef.current = editorStateToText(editor);
const rootElement = editor.getRootElement();
setScrollHeight(rootElement?.scrollHeight || 0);
startSts(() => {
......
......@@ -22,6 +22,7 @@ catalog:
'@modelcontextprotocol/sdk': ^1
'@fastgpt-sdk/storage': 0.6.15
'@fastgpt-sdk/logger': 0.1.2
'@fastgpt-sdk/otel': 0.1.0
'@types/lodash': ^4
'@types/react': ^18
'@types/react-dom': ^18
......@@ -44,6 +45,20 @@ catalog:
react: ^18
react-dom: ^18
react-i18next: 14.1.2
tsdown: ^0.21.0
tsdown: 0.21.4
typescript: ^5.9.3
zod: ^4
onlyBuiltDependencies:
- '@parcel/watcher'
- bufferutil
- canvas
- core-js
- esbuild
- mongodb-memory-server
- msgpackr-extract
- protobufjs
- puppeteer
- sharp
- utf-8-validate
- vue-demi
......@@ -50,12 +50,22 @@ HELPER_BOT_MODEL=qwen-max
# ==================== 日志配置 ====================
# 日志等级: trace | debug | info | warning | error | fatal
LOG_ENABLE_CONSOLE=true
LOG_CONSOLE_LEVEL=info
LOG_ENABLE_OTEL=false
LOG_CONSOLE_LEVEL=debug
LOG_ENABLE_OTEL=true
LOG_OTEL_LEVEL=info
LOG_OTEL_SERVICE_NAME=fastgpt-client
LOG_OTEL_URL=http://localhost:4318/v1/logs
# 指标
METRICS_ENABLE_OTEL=true
METRICS_OTEL_URL=http://localhost:4318/v1/metrics
METRICS_OTEL_SERVICE_NAME=fastgpt-client
# 追踪
TRACING_ENABLE_OTEL=true
TRACING_OTEL_URL=http://localhost:4318/v1/traces
TRACING_OTEL_SERVICE_NAME=fastgpt-client
# ==================== 对象存储 ====================
# 存储供应商;如果是 Sealos 的对象存储请填 aws-s3
STORAGE_VENDOR=minio
......
......@@ -45,6 +45,15 @@ ENV NODE_OPTIONS="--max-old-space-size=4096"
ENV NEXT_PUBLIC_BASE_URL=$base_url
RUN pnpm --filter=app build
# Remove build-time-only packages from standalone output before copying to runner.
# These are traced into standalone by mistake (rspack bindings, gnu platform binaries, etc.)
RUN rm -rf projects/app/.next/standalone/node_modules/.pnpm/@next+rspack-binding-*/ \
projects/app/.next/standalone/node_modules/.pnpm/@rspack+binding-*/ \
projects/app/.next/standalone/node_modules/.pnpm/next-rspack*/ \
projects/app/.next/standalone/node_modules/.pnpm/typescript@*/ \
projects/app/.next/standalone/node_modules/.pnpm/*-linux-x64-gnu@*/ \
projects/app/.next/standalone/node_modules/.pnpm/@img+sharp-libvips-linux-x64@*/
# --------- runner -----------
FROM node:20.14.0-alpine AS runner
WORKDIR /app
......@@ -74,18 +83,13 @@ COPY --from=builder --chown=nextjs:nodejs /app/projects/app/worker /app/projects
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
COPY --from=builder /app/projects/app/package.json ./package.json
# copy config
COPY ./projects/app/data/config.json /app/data/config.json
# copy test.mp3
COPY ./projects/app/data/test.mp3 /app/data/test.mp3
# copy GeoLite2-City.mmdb
COPY ./projects/app/data/GeoLite2-City.mmdb /app/data/GeoLite2-City.mmdb
RUN chown -R nextjs:nodejs /app/data
# Add tmp directory permission control
# copy config and data files (use --chown to avoid extra layer from chown)
COPY --chown=nextjs:nodejs ./projects/app/data/config.json /app/data/config.json
COPY --chown=nextjs:nodejs ./projects/app/data/test.mp3 /app/data/test.mp3
COPY --chown=nextjs:nodejs ./projects/app/data/GeoLite2-City.mmdb /app/data/GeoLite2-City.mmdb
ENV NODE_ENV=production
ENV NEXT_TELEMETRY_DISABLED=1
......
......@@ -62,6 +62,10 @@ const nextConfig: NextConfig = {
{
module: /bullmq[\\/]dist[\\/](cjs|esm)[\\/]classes[\\/]child-processor\.js$/,
message: /Critical dependency: the request of a dependency is an expression/
},
{
module: /@fastgpt-sdk[\\/]sandbox-adapter[\\/]/,
message: /Critical dependency/
}
];
......@@ -96,16 +100,14 @@ const nextConfig: NextConfig = {
}
if (isServer) {
(config.externals as string[]).push('@node-rs/jieba');
config.externals.push({
'@e2b/code-interpreter': 'commonjs @e2b/code-interpreter',
e2b: 'commonjs e2b'
'@node-rs/jieba': '@node-rs/jieba'
});
}
config.experiments = {
asyncWebAssembly: true,
layers: true
...config.experiments,
asyncWebAssembly: true
};
if (isDev && !isServer) {
......@@ -131,15 +133,14 @@ const nextConfig: NextConfig = {
return config;
},
transpilePackages: ['@modelcontextprotocol/sdk', 'ahooks', '@fastgpt-sdk/sandbox-adapter'],
transpilePackages: ['@modelcontextprotocol/sdk', 'ahooks'],
serverExternalPackages: [
'mongoose',
'pg',
'bullmq',
'@zilliz/milvus2-sdk-node',
'tiktoken',
'@opentelemetry/api-logs',
'chalk'
'@opentelemetry/api-logs'
],
// 优化大库的 barrel exports tree-shaking
experimental: {
......@@ -151,14 +152,37 @@ const nextConfig: NextConfig = {
'ahooks',
'framer-motion',
'@emotion/react',
'@emotion/styled'
'@emotion/styled',
'react-syntax-highlighter',
'recharts',
'@tanstack/react-query',
'react-hook-form',
'react-markdown'
],
// 按页面拆分 CSS chunk,减少首屏 CSS 体积
cssChunking: 'strict',
// 减少内存占用
memoryBasedWorkersCount: true
},
outputFileTracingRoot: path.join(__dirname, '../../')
outputFileTracingRoot: path.join(__dirname, '../../'),
// Exclude build-time-only packages from standalone output file tracing
outputFileTracingExcludes: {
'*': [
// Rspack bindings - only used in dev, not needed at runtime
'node_modules/@next/rspack-binding-*/**',
'node_modules/@rspack/binding-*/**',
'node_modules/next-rspack/**',
// GNU platform binaries - Alpine uses musl only
'node_modules/**/*-linux-x64-gnu*/**',
// typescript - build-time only
'node_modules/typescript/**',
// sharp libvips GNU variant (keep musl)
'node_modules/@img/sharp-libvips-linux-x64/**',
// bundle-analyzer - build-time only
'node_modules/@next/bundle-analyzer/**',
'node_modules/webpack-bundle-analyzer/**'
]
}
};
const configWithPluginsExceptWithRspack = withBundleAnalyzer(nextConfig);
......
{
"name": "app",
"version": "4.14.8.4",
"version": "4.14.9",
"private": false,
"scripts": {
"dev": "NODE_OPTIONS='--max-old-space-size=8192' npm run build:workers && next dev",
......
......@@ -80,7 +80,12 @@ const HumanContentCard = React.memo(
<Flex flexDirection={'column'} gap={4}>
{files.length > 0 && <FilesBlock files={files} />}
{text && (
<Box fontSize={'inherit'} color={'inherit'} whiteSpace={'pre-wrap'} wordBreak={'break-word'}>
<Box
fontSize={'inherit'}
color={'inherit'}
whiteSpace={'pre-wrap'}
wordBreak={'break-word'}
>
{text}
</Box>
)}
......
......@@ -46,7 +46,11 @@ const HumanItem = ({ chat }: { chat: UserChatItemType }) => {
>
<Flex flexDirection={'column'} gap={4}>
{files.length > 0 && <FileBlock files={files} />}
{text && <Box whiteSpace={'pre-wrap'} wordBreak={'break-word'}>{text}</Box>}
{text && (
<Box whiteSpace={'pre-wrap'} wordBreak={'break-word'}>
{text}
</Box>
)}
</Flex>
</Box>
<ChatAvatar type={ChatRoleEnum.Human} src={humanAvatar} />
......
......@@ -28,7 +28,8 @@ export const HelperBotContext = createContext<HelperBotContextType>({
taskObject: '',
selectedTools: [],
selectedDatasets: [],
fileUpload: false
fileUpload: false,
enableSandbox: false
},
onApply: function (e): void {
throw new Error('Function not implemented.');
......
......@@ -28,6 +28,8 @@ export async function register() {
{ initGeo },
{ instrumentationCheck },
{ getErrText },
{ configureMetrics },
{ configureTracing },
{ configureLogger, getLogger, LogCategories },
{ InitialErrorEnum }
] = await Promise.all([
......@@ -49,10 +51,14 @@ export async function register() {
import('@fastgpt/service/common/geo'),
import('@/service/common/system/health'),
import('@fastgpt/global/common/error/utils'),
import('@fastgpt/service/common/metrics'),
import('@fastgpt/service/common/tracing'),
import('@fastgpt/service/common/logger'),
import('@fastgpt/service/common/system/constants')
]);
await configureMetrics();
await configureTracing();
await configureLogger();
const logger = getLogger(LogCategories.SYSTEM);
logger.info('Starting system initialization...');
......
......@@ -78,6 +78,7 @@ const ChatTest = ({ appForm, setAppForm, setRenderEdit, form2WorkflowFn }: Props
selectedTools: appForm.selectedTools.map((tool) => tool.id),
selectedDatasets: appForm.dataset.datasets.map((dataset) => dataset.datasetId),
fileUpload: appForm.chatConfig.fileSelectConfig?.canSelectFile || false,
enableSandbox: appForm.aiSettings.useAgentSandbox || false,
modelConfig: {
model: appForm.aiSettings.model,
temperature: appForm.aiSettings.temperature,
......@@ -153,6 +154,7 @@ const ChatTest = ({ appForm, setAppForm, setRenderEdit, form2WorkflowFn }: Props
metadata={topAgentMetadata}
onApply={async (formData) => {
const fileUploadEnabled = !!formData.fileUploadEnabled;
const enableSandboxEnabled = !!formData.enableSandboxEnabled;
// Filter internal tools
const filteredToolIds = (formData.tools || []).filter(
......@@ -178,7 +180,8 @@ const ChatTest = ({ appForm, setAppForm, setRenderEdit, form2WorkflowFn }: Props
: prev.dataset,
aiSettings: {
...prev.aiSettings,
systemPrompt: formData.systemPrompt || prev.aiSettings.systemPrompt
systemPrompt: formData.systemPrompt || prev.aiSettings.systemPrompt,
useAgentSandbox: enableSandboxEnabled
},
chatConfig: {
...prev.chatConfig,
......
......@@ -99,7 +99,8 @@ const EditForm = ({
appForm.chatConfig.fileSelectConfig?.canSelectAudio ||
appForm.chatConfig.fileSelectConfig?.canSelectCustomFileExtension
),
hasSelectedDataset: (appForm.dataset.datasets?.length || 0) > 0
hasSelectedDataset: (appForm.dataset.datasets?.length || 0) > 0,
useAgentSandbox: !!appForm.aiSettings.useAgentSandbox
});
const {
......
......@@ -49,13 +49,15 @@ export const useSkillManager = ({
onUpdateOrAddTool,
onDeleteTool,
canUploadFile,
hasSelectedDataset
hasSelectedDataset,
useAgentSandbox
}: {
selectedTools: SelectedToolItemType[];
onDeleteTool: (id: string) => void;
onUpdateOrAddTool: (tool: SelectedToolItemType) => void;
canUploadFile: boolean;
hasSelectedDataset: boolean;
useAgentSandbox: boolean;
}) => {
const { t, i18n } = useTranslation();
const { toast } = useToast();
......@@ -109,6 +111,17 @@ export const useSkillManager = ({
});
}
const sandboxToolInfo = systemSubInfo[SubAppIds.sandboxTool];
if (sandboxToolInfo) {
apiTools.unshift({
id: SubAppIds.sandboxTool,
label: parseI18nString(sandboxToolInfo.name, i18n.language),
icon: sandboxToolInfo.avatar,
description: sandboxToolInfo.toolDescription,
canClick: true
});
}
return apiTools;
},
{
......@@ -324,8 +337,25 @@ export const useSkillManager = ({
});
}
// Merge sandbox tool
const sandboxToolInfo = systemSubInfo[SubAppIds.sandboxTool];
if (sandboxToolInfo) {
tools.push({
id: SubAppIds.sandboxTool,
pluginId: SubAppIds.sandboxTool,
name: parseI18nString(sandboxToolInfo.name, i18n.language),
avatar: sandboxToolInfo.avatar,
intro: sandboxToolInfo.toolDescription,
flowNodeType: FlowNodeTypeEnum.tool,
templateType: FlowNodeTemplateTypeEnum.tools,
inputs: [],
outputs: [],
configStatus: useAgentSandbox ? 'noConfig' : 'invalid'
});
}
return tools;
}, [selectedTools, canUploadFile, hasSelectedDataset, i18n.language]);
}, [selectedTools, canUploadFile, hasSelectedDataset, useAgentSandbox, i18n.language]);
const [configTool, setConfigTool] = useState<SelectedToolItemType>();
const onClickSkill = useCallback(
......
......@@ -223,12 +223,10 @@ const NodeCard = (props: Props) => {
// 1. MCP tool, HTTP tool set and system tool set do not have version
if (
isAppNode &&
(
node.toolConfig?.mcpToolSet ||
(node.toolConfig?.mcpToolSet ||
node.toolConfig?.mcpTool ||
node?.toolConfig?.httpToolSet ||
node?.toolConfig?.systemToolSet
)
node?.toolConfig?.systemToolSet)
)
return false;
// 2. Team app/System commercial plugin
......
# @fastgpt-sdk/otel
FastGPT 的统一 OpenTelemetry / observability SDK。
这个包的目标是作为未来的迁移目标,把现有的:
- `@fastgpt-sdk/logger`
- `@fastgpt-sdk/metrics`
- tracing 能力
收拢到一个统一入口里,但目前不强制迁移现有代码。
它现在是一个自包含包:
- 内部自带 logger 实现
- 内部自带 metrics 实现
- 内部自带 tracing 实现
- 不依赖 `@fastgpt-sdk/logger``@fastgpt-sdk/metrics`
同时支持两种使用方式:
- 统一入口:`@fastgpt-sdk/otel`
- 渐进迁移入口:`@fastgpt-sdk/otel/logger``@fastgpt-sdk/otel/metrics``@fastgpt-sdk/otel/tracing`
## 包含内容
- 内置 logger 能力
- 内置 metrics 能力
- 内置通用 tracing 能力
- 提供统一的 `configureOtel()` / `configureOtelFromEnv()` 入口
## 快速开始
```ts
import {
configureOtelFromEnv,
getLogger,
getMeter,
getTracer
} from '@fastgpt-sdk/otel';
await configureOtelFromEnv({
defaultServiceName: 'fastgpt-client'
});
const logger = getLogger(['system']);
const meter = getMeter('fastgpt-client');
const tracer = getTracer('fastgpt-client');
```
也可以渐进迁移:
```ts
import { configureLoggerFromEnv, getLogger } from '@fastgpt-sdk/otel/logger';
import { configureMetricsFromEnv, getMeter } from '@fastgpt-sdk/otel/metrics';
import { configureTracingFromEnv, getTracer } from '@fastgpt-sdk/otel/tracing';
```
## 迁移思路
未来可以分阶段迁移:
1. 先只把初始化入口从多个 SDK 收拢到 `@fastgpt-sdk/otel`
2. 再逐步把 import 从 `logger/metrics` 改成 `otel`
3. 最后按业务需要补 traces
## tracing 环境变量
- `TRACING_ENABLE_OTEL`
- `TRACING_OTEL_SERVICE_NAME`
- `TRACING_OTEL_URL`
- `TRACING_OTEL_SAMPLE_RATIO`
同时兼容标准 OTEL fallback:
- `OTEL_SERVICE_NAME`
- `OTEL_EXPORTER_OTLP_TRACES_ENDPOINT`
- `OTEL_EXPORTER_OTLP_ENDPOINT`
- `OTEL_TRACES_EXPORTER`
- `OTEL_TRACES_SAMPLER`
- `OTEL_TRACES_SAMPLER_ARG`
## 说明
- 这个包当前是“整理好的统一入口”,不是“已经迁移完成的替换方案”。
- 现有 `logger``metrics` 包仍然可继续独立使用,后续可以逐步迁移到这个包。
{
"name": "@fastgpt-sdk/otel",
"private": false,
"version": "0.1.0",
"description": "FastGPT SDK for OpenTelemetry observability",
"type": "module",
"main": "./dist/index.mjs",
"types": "./dist/index.d.mts",
"exports": {
".": {
"import": "./dist/index.mjs",
"types": "./dist/index.d.mts"
},
"./logger": {
"import": "./dist/logger-entry.mjs",
"types": "./dist/logger-entry.d.mts"
},
"./metrics": {
"import": "./dist/metrics-entry.mjs",
"types": "./dist/metrics-entry.d.mts"
},
"./tracing": {
"import": "./dist/tracing-entry.mjs",
"types": "./dist/tracing-entry.d.mts"
}
},
"files": [
"dist"
],
"scripts": {
"build": "tsdown",
"dev": "tsdown --watch",
"prepublishOnly": "pnpm build"
},
"keywords": [
"otel",
"opentelemetry",
"metrics",
"tracing",
"logging"
],
"author": "FastGPT",
"repository": {
"type": "git",
"url": "https://github.com/labring/FastGPT.git",
"directory": "FastGPT/sdk/otel"
},
"homepage": "https://github.com/labring/FastGPT",
"bugs": {
"url": "https://github.com/labring/FastGPT/issues"
},
"publishConfig": {
"access": "public"
},
"engines": {
"node": ">=20",
"pnpm": ">=9"
},
"packageManager": "pnpm@9.15.9",
"license": "Apache-2.0",
"dependencies": {
"@logtape/logtape": "^2",
"@logtape/pretty": "^2",
"@opentelemetry/api": "^1.9.0",
"@opentelemetry/api-logs": "^0.203.0",
"@opentelemetry/exporter-logs-otlp-http": "^0.203.0",
"@opentelemetry/exporter-metrics-otlp-http": "^0.203.0",
"@opentelemetry/exporter-trace-otlp-http": "^0.203.0",
"@opentelemetry/resources": "^2.0.1",
"@opentelemetry/sdk-logs": "^0.203.0",
"@opentelemetry/sdk-metrics": "^2.0.1",
"@opentelemetry/sdk-trace-base": "^2.0.1",
"@opentelemetry/sdk-trace-node": "^2.0.1",
"@opentelemetry/semantic-conventions": "^1.39.0"
},
"devDependencies": {
"@types/node": "catalog:",
"tsdown": "catalog:",
"typescript": "catalog:"
}
}
import { configureLogger, disposeLogger, getLogger } from './logger';
import type { LoggerConfigureOptions } from './logger';
import { configureMetrics, disposeMetrics, getMeter } from './metrics';
import type { MetricsConfigureOptions } from './metrics';
import { configureTracing, disposeTracing, getCurrentSpanContext, getTracer } from './tracing';
import type { TracingConfigureOptions } from './tracing';
import type { OtelConfigureOptions } from './types';
export async function configureOtel(options: OtelConfigureOptions = {}) {
await Promise.all([
configureLogger(options.logger ?? ({} satisfies LoggerConfigureOptions)),
configureMetrics(options.metrics ?? ({} satisfies MetricsConfigureOptions)),
configureTracing(options.tracing ?? ({} satisfies TracingConfigureOptions))
]);
}
export async function disposeOtel() {
await Promise.all([disposeLogger(), disposeMetrics(), disposeTracing()]);
}
export { getCurrentSpanContext, getLogger, getMeter, getTracer };
export type EnvValue = string | boolean | number | undefined;
export function parseBooleanEnv(value: EnvValue, defaultValue: boolean) {
if (typeof value === 'boolean') return value;
if (typeof value === 'number') return value !== 0;
if (typeof value !== 'string' || !value) return defaultValue;
const normalized = value.trim().toLowerCase();
if (['1', 'true', 'yes', 'on'].includes(normalized)) return true;
if (['0', 'false', 'no', 'off'].includes(normalized)) return false;
return defaultValue;
}
export function parseNumberEnv(value: EnvValue, defaultValue: number) {
if (typeof value === 'number' && Number.isFinite(value)) return value;
if (typeof value !== 'string') return defaultValue;
const parsed = Number(value);
return Number.isFinite(parsed) ? parsed : defaultValue;
}
export function parsePositiveNumberEnv(value: EnvValue, defaultValue: number) {
const parsed = parseNumberEnv(value, defaultValue);
return parsed > 0 ? parsed : defaultValue;
}
export function parseStringEnv(value: EnvValue): string | undefined {
if (typeof value !== 'string') return undefined;
const trimmed = value.trim();
return trimmed.length > 0 ? trimmed : undefined;
}
import {
configureLoggerFromEnv,
createLoggerOptionsFromEnv,
type LoggerConfigureFromEnvOptions,
type LoggerEnv
} from './logger';
import {
configureMetricsFromEnv,
createMetricsOptionsFromEnv,
type MetricsConfigureFromEnvOptions,
type MetricsEnv
} from './metrics';
import { configureTracingFromEnv, createTracingOptionsFromEnv } from './tracing';
import type { TracingConfigureFromEnvOptions, TracingEnv } from './tracing';
import type { OtelConfigureOptions } from './types';
type OtelEnv = LoggerEnv & MetricsEnv & TracingEnv;
export type OtelConfigureFromEnvOptions = {
env?: OtelEnv;
defaultServiceName?: string;
logger?: Omit<LoggerConfigureFromEnvOptions, 'env' | 'defaultServiceName'>;
metrics?: Omit<MetricsConfigureFromEnvOptions, 'env' | 'defaultServiceName'>;
tracing?: Omit<TracingConfigureFromEnvOptions, 'env' | 'defaultServiceName'>;
};
export function createOtelOptionsFromEnv(
options: OtelConfigureFromEnvOptions = {}
): OtelConfigureOptions {
const env = options.env ?? process.env;
return {
logger: createLoggerOptionsFromEnv({
env,
defaultServiceName: options.defaultServiceName,
...options.logger
}),
metrics: createMetricsOptionsFromEnv({
env,
defaultServiceName: options.defaultServiceName,
...options.metrics
}),
tracing: createTracingOptionsFromEnv({
env,
defaultServiceName: options.defaultServiceName,
...options.tracing
})
};
}
export async function configureOtelFromEnv(options: OtelConfigureFromEnvOptions = {}) {
const env = options.env ?? process.env;
await Promise.all([
configureLoggerFromEnv({
env,
defaultServiceName: options.defaultServiceName,
...options.logger
}),
configureMetricsFromEnv({
env,
defaultServiceName: options.defaultServiceName,
...options.metrics
}),
configureTracingFromEnv({
env,
defaultServiceName: options.defaultServiceName,
...options.tracing
})
]);
}
export { configureOtel, disposeOtel } from './client';
export { configureOtelFromEnv, createOtelOptionsFromEnv } from './env';
export * from './logger';
export * from './metrics';
export * from './tracing';
export type * from './types';
export * from './logger';
import { AsyncLocalStorage } from 'node:async_hooks';
import { configure, dispose, getLogger as getLogtapeLogger } from '@logtape/logtape';
import { createLoggers } from './loggers';
import { createSinks } from './sinks';
import type { LogCategory, LoggerConfigureOptions, LoggerContext } from './types';
let configured = false;
let configurePromise: Promise<void> | null = null;
let defaultCategory: LogCategory = ['system'];
export async function configureLogger(options: LoggerConfigureOptions = {}) {
if (configured) return;
if (configurePromise) return configurePromise;
configurePromise = (async () => {
defaultCategory = options.defaultCategory ?? defaultCategory;
const { sinks, composedSinks } = await createSinks({
console: options.console,
otel: options.otel,
sensitiveProperties: options.sensitiveProperties
});
const loggers = options.loggers ?? createLoggers({ composedSinks });
const contextLocalStorage =
options.contextLocalStorage ?? new AsyncLocalStorage<LoggerContext>();
await configure({
contextLocalStorage,
loggers,
sinks
});
configured = true;
})();
try {
await configurePromise;
} catch (error) {
configurePromise = null;
throw error;
}
}
export async function disposeLogger() {
if (configurePromise) {
try {
await configurePromise;
} catch {
configurePromise = null;
return;
}
}
if (!configured) return;
await dispose();
configured = false;
configurePromise = null;
}
export function getLogger(category: LogCategory = defaultCategory) {
const logger = getLogtapeLogger(category);
return new Proxy(logger, {
get(target, prop, receiver) {
const fn = Reflect.get(target, prop, receiver);
if (typeof fn !== 'function') return fn;
return (...args: unknown[]) => {
if (args.length === 0) return fn.call(target);
const [firstArg, secondArg] = args;
if (args.length === 1) {
return fn.call(target, firstArg);
}
if (typeof firstArg === 'string') {
if (
typeof secondArg === 'object' &&
secondArg &&
'verbose' in secondArg &&
typeof secondArg.verbose === 'boolean' &&
!secondArg.verbose
) {
const { verbose: _verbose, ...properties } = secondArg as Record<string, unknown> & {
verbose?: boolean;
};
return fn.call(target, firstArg, properties);
}
return fn.call(target, `${firstArg}: {*}`, secondArg);
}
if (typeof firstArg === 'object') {
return fn.call(target, firstArg);
}
return fn.apply(target, args);
};
}
});
}
import type { LogLevel } from '@logtape/logtape';
import { configureLogger } from './client';
import { parseBooleanEnv, parseStringEnv } from '../env-utils';
import type { LogCategory, LoggerConfigureOptions } from './types';
export type LoggerEnvValue = string | boolean | number | undefined;
export type LoggerEnv = Record<string, LoggerEnvValue>;
export type LoggerConfigureFromEnvOptions = {
env?: LoggerEnv;
defaultCategory?: LogCategory;
defaultServiceName?: string;
defaultLoggerName?: string;
defaultConsoleEnabled?: boolean;
defaultConsoleLevel?: LogLevel;
defaultOtelEnabled?: boolean;
defaultOtelLevel?: LogLevel;
defaultOtelUrl?: string;
sensitiveProperties?: readonly string[];
};
const logLevels = new Set<LogLevel>(['trace', 'debug', 'info', 'warning', 'error', 'fatal']);
function parseLogLevel(value: LoggerEnvValue, defaultValue: LogLevel): LogLevel {
if (typeof value !== 'string') return defaultValue;
return logLevels.has(value as LogLevel) ? (value as LogLevel) : defaultValue;
}
export function createLoggerOptionsFromEnv(
options: LoggerConfigureFromEnvOptions = {}
): LoggerConfigureOptions {
const env = options.env ?? process.env;
const defaultServiceName = options.defaultServiceName ?? 'app';
const serviceName = parseStringEnv(env.LOG_OTEL_SERVICE_NAME) ?? defaultServiceName;
const loggerName =
parseStringEnv(env.LOG_OTEL_LOGGER_NAME) ?? options.defaultLoggerName ?? serviceName;
return {
defaultCategory: options.defaultCategory,
console: {
enabled: parseBooleanEnv(env.LOG_ENABLE_CONSOLE, options.defaultConsoleEnabled ?? true),
level: parseLogLevel(env.LOG_CONSOLE_LEVEL, options.defaultConsoleLevel ?? 'trace')
},
otel: parseBooleanEnv(env.LOG_ENABLE_OTEL, options.defaultOtelEnabled ?? false)
? {
serviceName,
loggerName,
url:
parseStringEnv(env.LOG_OTEL_URL) ??
options.defaultOtelUrl ??
'http://localhost:4318/v1/logs',
level: parseLogLevel(env.LOG_OTEL_LEVEL, options.defaultOtelLevel ?? 'info')
}
: false,
sensitiveProperties: options.sensitiveProperties
};
}
export async function configureLoggerFromEnv(options: LoggerConfigureFromEnvOptions = {}) {
return configureLogger(createLoggerOptionsFromEnv(options));
}
import { SeverityNumber } from '@opentelemetry/api-logs';
export function mapLevelToSeverityNumber(level: string): number {
switch (level) {
case 'trace':
return SeverityNumber.TRACE;
case 'debug':
return SeverityNumber.DEBUG;
case 'info':
return SeverityNumber.INFO;
case 'warning':
return SeverityNumber.WARN;
case 'error':
return SeverityNumber.ERROR;
case 'fatal':
return SeverityNumber.FATAL;
default:
return SeverityNumber.UNSPECIFIED;
}
}
export { configureLogger, disposeLogger, getLogger } from './client';
export { withContext, withCategoryPrefix } from '@logtape/logtape';
export { getOpenTelemetrySink } from './otel';
export type {
BodyFormatter,
ExceptionAttributeMode,
ObjectRenderer,
OpenTelemetrySink,
OpenTelemetrySinkOptions
} from './otel';
export type {
ConsoleLoggerOptions,
LogCategory,
LoggerConfig,
LoggerConfigureOptions,
LoggerContext,
LoggerSinkId,
OtelLoggerOptions
} from './types';
export { configureLoggerFromEnv, createLoggerOptionsFromEnv } from './env';
export type { LoggerConfigureFromEnvOptions, LoggerEnv } from './env';
import type { LogTapeConfig, LoggerSinkId } from './types';
type LoggerConfig = LogTapeConfig['loggers'];
type CreateLoggersOptions = {
composedSinks: LoggerSinkId[];
};
export function createLoggers({ composedSinks }: CreateLoggersOptions): LoggerConfig {
const metaSinks: LoggerSinkId[] = composedSinks.includes('console') ? ['console'] : composedSinks;
return [
{
category: [],
lowestLevel: 'trace',
sinks: composedSinks
},
...(metaSinks.length === 0
? []
: [
{
category: ['logtape', 'meta'],
lowestLevel: 'fatal' as const,
parentSinks: 'override' as const,
sinks: metaSinks
}
])
];
}
import type { LogLevel, LogRecord } from '@logtape/logtape';
import { getConsoleSink, withFilter } from '@logtape/logtape';
import { getPrettyFormatter } from '@logtape/pretty';
import { mapLevelToSeverityNumber } from './helpers';
import { getOpenTelemetrySink } from './otel';
import type {
ConsoleLoggerOptions,
LogTapeConfig,
LoggerConfigureOptions,
LoggerSinkId,
OtelLoggerOptions
} from './types';
type SinkConfig = LogTapeConfig<string>['sinks'];
type CreateSinksOptions = Pick<LoggerConfigureOptions, 'console' | 'otel' | 'sensitiveProperties'>;
type CreateSinksResult = {
sinks: SinkConfig;
composedSinks: LoggerSinkId[];
};
const defaultConsoleOptions: Required<ConsoleLoggerOptions> = {
enabled: true,
level: 'trace'
};
const defaultOtelOptions = {
enabled: false,
level: 'info' as LogLevel
};
function normalizeConsoleOptions(
options?: boolean | ConsoleLoggerOptions
): Required<ConsoleLoggerOptions> {
if (typeof options === 'boolean') {
return {
...defaultConsoleOptions,
enabled: options
};
}
return {
enabled: options?.enabled ?? defaultConsoleOptions.enabled,
level: options?.level ?? defaultConsoleOptions.level
};
}
function normalizeOtelOptions(options?: false | OtelLoggerOptions) {
if (!options) {
return {
...defaultOtelOptions,
serviceName: undefined,
url: undefined,
loggerName: undefined
};
}
return {
enabled: options.enabled ?? true,
level: options.level ?? defaultOtelOptions.level,
serviceName: options.serviceName,
url: options.url,
loggerName: options.loggerName ?? options.serviceName
};
}
function pad(value: number) {
return value.toString().padStart(2, '0');
}
function formatTimestamp(timestamp: number | Date) {
const date = timestamp instanceof Date ? timestamp : new Date(timestamp);
return `${date.getFullYear()}-${pad(date.getMonth() + 1)}-${pad(date.getDate())} ${pad(
date.getHours()
)}:${pad(date.getMinutes())}:${pad(date.getSeconds())}`;
}
export async function createSinks(options: CreateSinksOptions): Promise<CreateSinksResult> {
const consoleOptions = normalizeConsoleOptions(options.console);
const otelOptions = normalizeOtelOptions(options.otel);
const sensitiveProperties = options.sensitiveProperties ?? [];
const sinkConfig = {
bufferSize: 8192,
flushInterval: 5000,
nonBlocking: true,
lazy: true
} as const;
const sinks: SinkConfig = {};
const composedSinks: LoggerSinkId[] = [];
const levelFilter = (record: LogRecord, level: LogLevel) => {
return mapLevelToSeverityNumber(record.level) >= mapLevelToSeverityNumber(level);
};
if (consoleOptions.enabled) {
sinks.console = withFilter(
getConsoleSink({
...sinkConfig,
formatter: getPrettyFormatter({
icons: false,
level: 'ABBR',
wordWrap: false,
messageColor: null,
categoryColor: null,
timestampColor: null,
levelStyle: 'reset',
messageStyle: 'reset',
categoryStyle: 'reset',
timestampStyle: 'reset',
categorySeparator: ':',
timestamp: formatTimestamp,
inspectOptions: { depth: 5 }
})
}),
(record) => levelFilter(record, consoleOptions.level)
);
composedSinks.push('console');
}
if (otelOptions.enabled) {
if (!otelOptions.serviceName) {
throw new Error('`otel.serviceName` is required when OpenTelemetry logging is enabled');
}
sinks.otel = withFilter(
getOpenTelemetrySink({
serviceName: otelOptions.serviceName,
loggerName: otelOptions.loggerName,
otlpExporterConfig: otelOptions.url ? { url: otelOptions.url } : undefined
}),
(record) => {
const properties = record.properties ?? {};
return (
levelFilter(record, otelOptions.level) &&
!sensitiveProperties.some((property) => property in properties)
);
}
);
composedSinks.push('otel');
}
return { sinks, composedSinks };
}
import type { AsyncLocalStorage } from 'node:async_hooks';
import type { Config, LogLevel } from '@logtape/logtape';
export type LogCategory = readonly string[];
export type LoggerContext = Record<string, unknown>;
export type LoggerSinkId = 'console' | 'otel';
type FilterId = string;
export type LogTapeConfig<S extends string = LoggerSinkId, F extends string = FilterId> = Config<
S,
F
>;
export type LoggerConfig = LogTapeConfig['loggers'];
export type ConsoleLoggerOptions = {
enabled?: boolean;
level?: LogLevel;
};
export type OtelLoggerOptions = {
enabled?: boolean;
level?: LogLevel;
serviceName: string;
url?: string;
loggerName?: string;
};
export type LoggerConfigureOptions = {
console?: boolean | ConsoleLoggerOptions;
otel?: false | OtelLoggerOptions;
contextLocalStorage?: AsyncLocalStorage<LoggerContext>;
loggers?: LoggerConfig;
sensitiveProperties?: readonly string[];
defaultCategory?: LogCategory;
};
export * from './metrics';
import { metrics } from '@opentelemetry/api';
import { OTLPMetricExporter } from '@opentelemetry/exporter-metrics-otlp-http';
import { defaultResource, resourceFromAttributes } from '@opentelemetry/resources';
import { MeterProvider, PeriodicExportingMetricReader } from '@opentelemetry/sdk-metrics';
import { ATTR_SERVICE_NAME } from '@opentelemetry/semantic-conventions';
import type { MetricsConfigureOptions, MetricsOptions } from './types';
type OtlpMetricExporterConfig = ConstructorParameters<typeof OTLPMetricExporter>[0];
let configured = false;
let configurePromise: Promise<void> | null = null;
let meterProvider: MeterProvider | null = null;
let defaultMeterName = 'fastgpt';
let defaultMeterVersion: string | undefined;
function getEnvironmentVariable(name: string): string | undefined {
return process.env[name];
}
function hasOtlpEndpoint(config?: OtlpMetricExporterConfig): boolean {
if (config?.url) return true;
if (getEnvironmentVariable('OTEL_EXPORTER_OTLP_METRICS_ENDPOINT')) return true;
if (getEnvironmentVariable('OTEL_EXPORTER_OTLP_ENDPOINT')) return true;
return false;
}
function normalizeOtlpMetricsUrl(url: string) {
const trimmed = url.trim();
if (!trimmed) return trimmed;
if (trimmed.endsWith('/v1/metrics')) return trimmed;
return `${trimmed.replace(/\/+$/, '')}/v1/metrics`;
}
function resolveOtlpMetricsUrl(config?: OtlpMetricExporterConfig) {
if (config?.url) return config.url;
const metricsEndpoint = getEnvironmentVariable('OTEL_EXPORTER_OTLP_METRICS_ENDPOINT');
if (metricsEndpoint) return metricsEndpoint;
const endpoint = getEnvironmentVariable('OTEL_EXPORTER_OTLP_ENDPOINT');
if (endpoint) return normalizeOtlpMetricsUrl(endpoint);
return undefined;
}
function normalizeMetricsOptions(options?: false | MetricsOptions) {
if (options === false) {
return {
enabled: false,
exportIntervalMillis: 15000
};
}
return {
enabled: options?.enabled ?? false,
serviceName: options?.serviceName,
exportIntervalMillis: options?.exportIntervalMillis ?? 15000,
otlpExporterConfig: {
url: options?.url,
headers: options?.headers
} satisfies OtlpMetricExporterConfig,
additionalResource: options?.additionalResource ?? null
};
}
export async function configureMetrics(options: MetricsConfigureOptions = {}) {
if (configured) return;
if (configurePromise) return configurePromise;
configurePromise = (async () => {
const metricsOptions = normalizeMetricsOptions(options.metrics);
defaultMeterName = options.defaultMeterName ?? defaultMeterName;
defaultMeterVersion = options.defaultMeterVersion ?? defaultMeterVersion;
const resource = defaultResource().merge(
resourceFromAttributes({
[ATTR_SERVICE_NAME]:
metricsOptions.serviceName ??
getEnvironmentVariable('OTEL_SERVICE_NAME') ??
defaultMeterName
}).merge(metricsOptions.additionalResource ?? null)
);
const readers: PeriodicExportingMetricReader[] = [];
if (metricsOptions.enabled && hasOtlpEndpoint(metricsOptions.otlpExporterConfig)) {
const exporter = new OTLPMetricExporter({
...metricsOptions.otlpExporterConfig,
url: resolveOtlpMetricsUrl(metricsOptions.otlpExporterConfig)
});
readers.push(
new PeriodicExportingMetricReader({
exporter,
exportIntervalMillis: metricsOptions.exportIntervalMillis
})
);
}
meterProvider = new MeterProvider({
resource,
readers
});
metrics.setGlobalMeterProvider(meterProvider);
configured = true;
})();
try {
await configurePromise;
} catch (error) {
configurePromise = null;
throw error;
}
}
export async function disposeMetrics() {
if (configurePromise) {
try {
await configurePromise;
} catch {
configurePromise = null;
return;
}
}
if (!configured || !meterProvider) return;
await meterProvider.shutdown();
configured = false;
configurePromise = null;
meterProvider = null;
}
export function getMeter(name = defaultMeterName, version = defaultMeterVersion) {
return metrics.getMeter(name, version);
}
import { configureMetrics } from './client';
import { parseBooleanEnv, parsePositiveNumberEnv, parseStringEnv } from '../env-utils';
import type { MetricsConfigureOptions } from './types';
export type MetricsEnvValue = string | boolean | number | undefined;
export type MetricsEnv = Record<string, MetricsEnvValue>;
export type MetricsConfigureFromEnvOptions = {
env?: MetricsEnv;
defaultServiceName?: string;
defaultMeterName?: string;
defaultMetricsEnabled?: boolean;
defaultMetricsUrl?: string;
defaultExportIntervalMillis?: number;
};
export function createMetricsOptionsFromEnv(
options: MetricsConfigureFromEnvOptions = {}
): MetricsConfigureOptions {
const env = options.env ?? process.env;
const metricsExporter = parseStringEnv(env.OTEL_METRICS_EXPORTER)?.toLowerCase();
const enabled = parseBooleanEnv(
env.METRICS_ENABLE_OTEL,
metricsExporter === 'otlp' || options.defaultMetricsEnabled === true
);
return {
defaultMeterName: options.defaultMeterName ?? options.defaultServiceName ?? 'fastgpt',
metrics: enabled
? {
enabled: true,
serviceName:
parseStringEnv(env.METRICS_OTEL_SERVICE_NAME) ??
parseStringEnv(env.OTEL_SERVICE_NAME) ??
options.defaultServiceName,
url:
parseStringEnv(env.METRICS_OTEL_URL) ??
parseStringEnv(env.OTEL_EXPORTER_OTLP_METRICS_ENDPOINT) ??
options.defaultMetricsUrl,
exportIntervalMillis: parsePositiveNumberEnv(
env.METRICS_EXPORT_INTERVAL ?? env.OTEL_METRIC_EXPORT_INTERVAL,
options.defaultExportIntervalMillis ?? 15000
)
}
: false
};
}
export async function configureMetricsFromEnv(options: MetricsConfigureFromEnvOptions = {}) {
return configureMetrics(createMetricsOptionsFromEnv(options));
}
export { configureMetrics, disposeMetrics, getMeter } from './client';
export { configureMetricsFromEnv, createMetricsOptionsFromEnv } from './env';
export type { MetricsConfigureOptions, MetricsOptions, MetricAttributes } from './types';
export type { MetricsConfigureFromEnvOptions, MetricsEnv } from './env';
import type { Resource } from '@opentelemetry/resources';
export type MetricsOptions = {
enabled?: boolean;
serviceName?: string;
url?: string;
headers?: Record<string, string>;
exportIntervalMillis?: number;
additionalResource?: Resource | null;
};
export type MetricsConfigureOptions = {
defaultMeterName?: string;
defaultMeterVersion?: string;
metrics?: false | MetricsOptions;
};
export type MetricAttributeValue = string | number | boolean;
export type MetricAttributes = Record<string, MetricAttributeValue>;
export { configureTracing, disposeTracing, getCurrentSpanContext, getTracer } from './tracing';
export { configureTracingFromEnv, createTracingOptionsFromEnv } from './tracing';
export type * from './tracing';
import { trace } from '@opentelemetry/api';
import { OTLPTraceExporter } from '@opentelemetry/exporter-trace-otlp-http';
import { defaultResource, resourceFromAttributes } from '@opentelemetry/resources';
import {
BatchSpanProcessor,
ParentBasedSampler,
TraceIdRatioBasedSampler
} from '@opentelemetry/sdk-trace-base';
import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node';
import { ATTR_SERVICE_NAME } from '@opentelemetry/semantic-conventions';
import type { TracingConfigureOptions, TracingOptions } from './types';
type OtlpTraceExporterConfig = ConstructorParameters<typeof OTLPTraceExporter>[0];
let configured = false;
let configurePromise: Promise<void> | null = null;
let tracerProvider: NodeTracerProvider | null = null;
let defaultTracerName = 'fastgpt';
let defaultTracerVersion: string | undefined;
function getEnvironmentVariable(name: string): string | undefined {
return process.env[name];
}
function hasOtlpEndpoint(config?: OtlpTraceExporterConfig): boolean {
if (config?.url) return true;
if (getEnvironmentVariable('OTEL_EXPORTER_OTLP_TRACES_ENDPOINT')) return true;
if (getEnvironmentVariable('OTEL_EXPORTER_OTLP_ENDPOINT')) return true;
return false;
}
function normalizeOtlpTracesUrl(url: string) {
const trimmed = url.trim();
if (!trimmed) return trimmed;
if (trimmed.endsWith('/v1/traces')) return trimmed;
return `${trimmed.replace(/\/+$/, '')}/v1/traces`;
}
function resolveOtlpTracesUrl(config?: OtlpTraceExporterConfig) {
if (config?.url) return config.url;
const tracesEndpoint = getEnvironmentVariable('OTEL_EXPORTER_OTLP_TRACES_ENDPOINT');
if (tracesEndpoint) return tracesEndpoint;
const endpoint = getEnvironmentVariable('OTEL_EXPORTER_OTLP_ENDPOINT');
if (endpoint) return normalizeOtlpTracesUrl(endpoint);
return undefined;
}
function normalizeSampleRatio(value: number | undefined, defaultValue: number) {
if (typeof value !== 'number' || !Number.isFinite(value)) return defaultValue;
return Math.max(0, Math.min(1, value));
}
function normalizeTracingOptions(options?: false | TracingOptions) {
if (options === false) {
return {
enabled: false,
sampleRatio: 1
};
}
return {
enabled: options?.enabled ?? false,
serviceName: options?.serviceName,
sampleRatio: normalizeSampleRatio(options?.sampleRatio, 1),
otlpExporterConfig: {
url: options?.url,
headers: options?.headers
} satisfies OtlpTraceExporterConfig,
additionalResource: options?.additionalResource ?? null
};
}
export async function configureTracing(options: TracingConfigureOptions = {}) {
if (configured) return;
if (configurePromise) return configurePromise;
configurePromise = (async () => {
const tracingOptions = normalizeTracingOptions(options.tracing);
defaultTracerName = options.defaultTracerName ?? defaultTracerName;
defaultTracerVersion = options.defaultTracerVersion ?? defaultTracerVersion;
if (!tracingOptions.enabled) {
configured = true;
return;
}
const resource = defaultResource().merge(
resourceFromAttributes({
[ATTR_SERVICE_NAME]:
tracingOptions.serviceName ??
getEnvironmentVariable('OTEL_SERVICE_NAME') ??
defaultTracerName
}).merge(tracingOptions.additionalResource ?? null)
);
const spanProcessors = [];
if (hasOtlpEndpoint(tracingOptions.otlpExporterConfig)) {
const exporter = new OTLPTraceExporter({
...tracingOptions.otlpExporterConfig,
url: resolveOtlpTracesUrl(tracingOptions.otlpExporterConfig)
});
spanProcessors.push(new BatchSpanProcessor(exporter));
}
tracerProvider = new NodeTracerProvider({
resource,
sampler: new ParentBasedSampler({
root: new TraceIdRatioBasedSampler(tracingOptions.sampleRatio)
}),
spanProcessors
});
tracerProvider.register();
configured = true;
})();
try {
await configurePromise;
} catch (error) {
configurePromise = null;
throw error;
}
}
export async function disposeTracing() {
if (configurePromise) {
try {
await configurePromise;
} catch {
configurePromise = null;
return;
}
}
if (!configured) return;
if (!tracerProvider) {
configured = false;
configurePromise = null;
return;
}
await tracerProvider.shutdown();
configured = false;
configurePromise = null;
tracerProvider = null;
}
export function getTracer(name = defaultTracerName, version = defaultTracerVersion) {
return trace.getTracer(name, version);
}
export function getCurrentSpanContext() {
return trace.getActiveSpan()?.spanContext();
}
import { configureTracing } from './client';
import { parseBooleanEnv, parseNumberEnv, parseStringEnv } from '../env-utils';
import type { TracingConfigureOptions } from './types';
export type TracingEnvValue = string | boolean | number | undefined;
export type TracingEnv = Record<string, TracingEnvValue>;
export type TracingConfigureFromEnvOptions = {
env?: TracingEnv;
defaultServiceName?: string;
defaultTracerName?: string;
defaultTracingEnabled?: boolean;
defaultTracingUrl?: string;
defaultSampleRatio?: number;
};
function normalizeSampleRatio(value: number, defaultValue: number) {
if (!Number.isFinite(value)) return defaultValue;
return Math.max(0, Math.min(1, value));
}
function getSampleRatioFromStandardEnv(env: TracingEnv, defaultValue: number): number {
const sampler = parseStringEnv(env.OTEL_TRACES_SAMPLER)?.toLowerCase();
const samplerArg = normalizeSampleRatio(
parseNumberEnv(env.OTEL_TRACES_SAMPLER_ARG, defaultValue),
defaultValue
);
if (sampler === 'always_off' || sampler === 'parentbased_always_off') return 0;
if (sampler === 'always_on' || sampler === 'parentbased_always_on') return 1;
if (sampler === 'traceidratio' || sampler === 'parentbased_traceidratio') {
return samplerArg;
}
return defaultValue;
}
export function createTracingOptionsFromEnv(
options: TracingConfigureFromEnvOptions = {}
): TracingConfigureOptions {
const env = options.env ?? process.env;
const tracesExporter = parseStringEnv(env.OTEL_TRACES_EXPORTER)?.toLowerCase();
const enabled = parseBooleanEnv(
env.TRACING_ENABLE_OTEL,
tracesExporter === 'otlp' || options.defaultTracingEnabled === true
);
const defaultSampleRatio = normalizeSampleRatio(options.defaultSampleRatio ?? 1, 1);
return {
defaultTracerName: options.defaultTracerName ?? options.defaultServiceName ?? 'fastgpt',
tracing: enabled
? {
enabled: true,
serviceName:
parseStringEnv(env.TRACING_OTEL_SERVICE_NAME) ??
parseStringEnv(env.OTEL_SERVICE_NAME) ??
options.defaultServiceName,
url:
parseStringEnv(env.TRACING_OTEL_URL) ??
parseStringEnv(env.OTEL_EXPORTER_OTLP_TRACES_ENDPOINT) ??
options.defaultTracingUrl,
sampleRatio: normalizeSampleRatio(
parseNumberEnv(env.TRACING_OTEL_SAMPLE_RATIO, NaN),
getSampleRatioFromStandardEnv(env, defaultSampleRatio)
)
}
: false
};
}
export async function configureTracingFromEnv(options: TracingConfigureFromEnvOptions = {}) {
return configureTracing(createTracingOptionsFromEnv(options));
}
export { configureTracing, disposeTracing, getCurrentSpanContext, getTracer } from './client';
export { configureTracingFromEnv, createTracingOptionsFromEnv } from './env';
export type * from './types';
export type * from './env';
import type { Resource } from '@opentelemetry/resources';
export type TracingOptions = {
enabled?: boolean;
serviceName?: string;
url?: string;
headers?: Record<string, string>;
sampleRatio?: number;
additionalResource?: Resource | null;
};
export type TracingConfigureOptions = {
defaultTracerName?: string;
defaultTracerVersion?: string;
tracing?: false | TracingOptions;
};
export type TraceAttributeValue = string | number | boolean;
export type TraceAttributes = Record<string, TraceAttributeValue>;
import type { LoggerConfigureOptions } from './logger';
import type { MetricsConfigureOptions } from './metrics';
import type { TracingConfigureOptions } from './tracing';
export type OtelConfigureOptions = {
logger?: LoggerConfigureOptions;
metrics?: MetricsConfigureOptions;
tracing?: TracingConfigureOptions;
};
{
"compilerOptions": {
"module": "esnext",
"target": "es2022",
"moduleResolution": "bundler",
"sourceMap": true,
"declaration": true,
"declarationMap": true,
"strict": true,
"verbatimModuleSyntax": true,
"isolatedModules": true,
"noUncheckedSideEffectImports": true,
"moduleDetection": "force",
"skipLibCheck": true,
"noEmit": true
}
}
import { defineConfig } from 'tsdown';
export default defineConfig({
entry: ['src/index.ts', 'src/logger-entry.ts', 'src/metrics-entry.ts', 'src/tracing-entry.ts'],
format: 'esm',
dts: {
enabled: true,
sourcemap: false
}
});
import { describe, it, expect } from 'vitest';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import { WorkflowQueue } from '@fastgpt/service/core/workflow/dispatch/index';
import { createNode, createEdge, setEdgeStatus } from '../../utils';
describe('场景21: 混合边状态 - 部分 active、部分 waiting、部分 skipped', () => {
/**
* 工作流结构:
*
* ┌──→ A ──┐
* start ┤ ├──→ D
* ├──→ B ──┤
* └──→ C ──┘
*
* 预期分组:
* - D: 组1[A→D, B→D, C→D] (并行汇聚,所有边在同一组)
*
* 测试场景:
* 1. A active, B waiting, C skipped → D 应该等待
* 2. A active, B active, C skipped → D 应该运行
* 3. A skipped, B skipped, C skipped → D 应该跳过
* 4. A waiting, B waiting, C waiting → D 应该等待
*/
const nodes = [
createNode('start', FlowNodeTypeEnum.workflowStart),
createNode('A', FlowNodeTypeEnum.chatNode),
createNode('B', FlowNodeTypeEnum.chatNode),
createNode('C', FlowNodeTypeEnum.chatNode),
createNode('D', FlowNodeTypeEnum.chatNode)
];
const edges = [
createEdge('start', 'A'),
createEdge('start', 'B'),
createEdge('start', 'C'),
createEdge('A', 'D'),
createEdge('B', 'D'),
createEdge('C', 'D')
];
const edgeIndex = WorkflowQueue.buildEdgeIndex({ runtimeEdges: edges });
const edgeGroupsMap = WorkflowQueue.buildNodeEdgeGroupsMap({
runtimeNodes: nodes,
edgeIndex
});
it('D 节点应该只有 1 组(并行汇聚)', () => {
const groups = edgeGroupsMap.get('D') || [];
expect(groups.length).toBe(1);
});
it('A active, B waiting, C skipped → D 应该等待', () => {
setEdgeStatus(edges, 'A', 'D', 'active');
setEdgeStatus(edges, 'B', 'D', 'waiting');
setEdgeStatus(edges, 'C', 'D', 'skipped');
const statusD = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'D')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(statusD).toBe('wait');
});
it('A active, B active, C skipped → D 应该运行', () => {
setEdgeStatus(edges, 'A', 'D', 'active');
setEdgeStatus(edges, 'B', 'D', 'active');
setEdgeStatus(edges, 'C', 'D', 'skipped');
const statusD = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'D')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(statusD).toBe('run');
});
it('A skipped, B skipped, C skipped → D 应该跳过', () => {
setEdgeStatus(edges, 'A', 'D', 'skipped');
setEdgeStatus(edges, 'B', 'D', 'skipped');
setEdgeStatus(edges, 'C', 'D', 'skipped');
const statusD = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'D')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(statusD).toBe('skip');
});
it('A waiting, B waiting, C waiting → D 应该等待', () => {
setEdgeStatus(edges, 'A', 'D', 'waiting');
setEdgeStatus(edges, 'B', 'D', 'waiting');
setEdgeStatus(edges, 'C', 'D', 'waiting');
const statusD = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'D')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(statusD).toBe('wait');
});
});
describe('场景22: 孤立节点和终止节点', () => {
/**
* 工作流结构:
*
* start → A → B
*
* C (孤立节点,没有输入边)
*
* 测试场景:
* 1. B 节点没有输出边(终止节点)
* 2. C 节点没有输入边(孤立节点)- 实际上没有输入边的节点会被视为可以运行
*/
const nodes = [
createNode('start', FlowNodeTypeEnum.workflowStart),
createNode('A', FlowNodeTypeEnum.chatNode),
createNode('B', FlowNodeTypeEnum.chatNode),
createNode('C', FlowNodeTypeEnum.chatNode)
];
const edges = [createEdge('start', 'A'), createEdge('A', 'B')];
const edgeIndex = WorkflowQueue.buildEdgeIndex({ runtimeEdges: edges });
const edgeGroupsMap = WorkflowQueue.buildNodeEdgeGroupsMap({
runtimeNodes: nodes,
edgeIndex
});
it('B 节点(终止节点)应该能正常运行', () => {
setEdgeStatus(edges, 'A', 'B', 'active');
const statusB = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'B')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(statusB).toBe('run');
});
it('C 节点(孤立节点)没有输入边分组,应该返回 run', () => {
const groups = edgeGroupsMap.get('C') || [];
expect(groups.length).toBe(0);
const statusC = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'C')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
// 没有输入边的节点,getNodeRunStatus 会返回 'run'
expect(statusC).toBe('run');
});
});
describe('场景23: userSelect 节点的多选项分支', () => {
/**
* 工作流结构:
*
* ┌──option1──→ A ──┐
* start → Select ┤──option2──→ B ──┼──→ End
* └──option3──→ C ──┘
*
* 预期分组:
* - A: 组1[Select→A (option1 handle)]
* - B: 组1[Select→B (option2 handle)]
* - C: 组1[Select→C (option3 handle)]
* - End: 组1[A→End, B→End, C→End] (并行汇聚,所有边在同一组)
*
* 测试场景:
* 1. 选择 option1:A 应该运行,B 和 C 应该跳过
* 2. 选择 option2:B 应该运行,A 和 C 应该跳过
* 3. 选择 option3:C 应该运行,A 和 B 应该跳过
*/
const nodes = [
createNode('start', FlowNodeTypeEnum.workflowStart),
createNode('Select', FlowNodeTypeEnum.userSelect),
createNode('A', FlowNodeTypeEnum.chatNode),
createNode('B', FlowNodeTypeEnum.chatNode),
createNode('C', FlowNodeTypeEnum.chatNode),
createNode('End', FlowNodeTypeEnum.chatNode)
];
const edges = [
createEdge('start', 'Select'),
createEdge('Select', 'A', 'waiting', 'Select-source-option1', 'A-target-left'),
createEdge('Select', 'B', 'waiting', 'Select-source-option2', 'B-target-left'),
createEdge('Select', 'C', 'waiting', 'Select-source-option3', 'C-target-left'),
createEdge('A', 'End'),
createEdge('B', 'End'),
createEdge('C', 'End')
];
const edgeIndex = WorkflowQueue.buildEdgeIndex({ runtimeEdges: edges });
const edgeGroupsMap = WorkflowQueue.buildNodeEdgeGroupsMap({
runtimeNodes: nodes,
edgeIndex
});
it('End 节点应该只有 1 组(并行汇聚)', () => {
const groups = edgeGroupsMap.get('End') || [];
expect(groups.length).toBe(1);
});
it('选择 option1:A 应该运行,B 和 C 应该跳过', () => {
setEdgeStatus(edges, 'Select', 'A', 'active');
setEdgeStatus(edges, 'Select', 'B', 'skipped');
setEdgeStatus(edges, 'Select', 'C', 'skipped');
expect(
WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'A')!,
nodeEdgeGroupsMap: edgeGroupsMap
})
).toBe('run');
expect(
WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'B')!,
nodeEdgeGroupsMap: edgeGroupsMap
})
).toBe('skip');
expect(
WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'C')!,
nodeEdgeGroupsMap: edgeGroupsMap
})
).toBe('skip');
});
it('选择 option2:B 应该运行,A 和 C 应该跳过', () => {
setEdgeStatus(edges, 'Select', 'A', 'skipped');
setEdgeStatus(edges, 'Select', 'B', 'active');
setEdgeStatus(edges, 'Select', 'C', 'skipped');
expect(
WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'A')!,
nodeEdgeGroupsMap: edgeGroupsMap
})
).toBe('skip');
expect(
WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'B')!,
nodeEdgeGroupsMap: edgeGroupsMap
})
).toBe('run');
expect(
WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'C')!,
nodeEdgeGroupsMap: edgeGroupsMap
})
).toBe('skip');
});
it('A 执行完成后,End 应该运行', () => {
setEdgeStatus(edges, 'Select', 'A', 'active');
setEdgeStatus(edges, 'Select', 'B', 'skipped');
setEdgeStatus(edges, 'Select', 'C', 'skipped');
setEdgeStatus(edges, 'A', 'End', 'active');
setEdgeStatus(edges, 'B', 'End', 'skipped');
setEdgeStatus(edges, 'C', 'End', 'skipped');
const statusEnd = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'End')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(statusEnd).toBe('run');
});
});
import { describe, it, expect } from 'vitest';
import { FlowNodeTypeEnum } from '@fastgpt/global/core/workflow/node/constant';
import { WorkflowQueue } from '@fastgpt/service/core/workflow/dispatch/index';
import { createNode, createEdge, setEdgeStatus } from '../../utils';
describe('1: 医疗记录工作流 - 非对称分支汇聚 + 中间节点循环', () => {
/**
* 工作流结构(来源于真实用户工作流):
*
* IF → code → newArr → getFirst → updateArr → updateCur1 ──┐
* start → ifElse1 ↑ ├──→ AI → ifElse2
* ELSE ──── updateCur2 ───┼────────────────────────────────┘ │
* │ IF → updateHistory ─┘
* │ ELSE → reply [退出]
* └─────────────────────────────────────────────────────────────
*
* 关键特点:
* 1. ifElse1 的 IF 分支经过长链路(code→newArr→getFirst→updateArr→updateCur1),ELSE 分支走短链路(updateCur2)
* 2. 两条分支汇聚到 AI 节点(非对称汇聚)
* 3. 循环目标是中间节点 getFirst,而非初始 ifElse1(非对称循环)
* 4. getFirst 有两个来源组:[newArr→getFirst](初始 IF 路径)和 [updateHistory→getFirst](循环路径)
* 5. AI 有两个来源组:[updateCur1→AI](IF 路径)和 [updateCur2→AI](ELSE 路径)
*
* 预期分组:
* - getFirst: 组1[newArr→getFirst], 组2[updateHistory→getFirst]
* - AI: 组1[updateCur1→AI], 组2[updateCur2→AI]
*
* 测试场景:
* 1. IF 分支首次执行:newArr→getFirst active → getFirst 应该运行
* 2. IF 分支首次执行:updateCur1���AI active, updateCur2→AI skipped → AI 应该运行
* 3. ELSE 分支执行:updateCur2→AI active, updateCur1→AI skipped → AI 应该运行
* 4. ELSE 分支执行:newArr→getFirst skipped → getFirst 应该跳过
* 5. 循环迭代:updateHistory→getFirst active, newArr→getFirst skipped → getFirst 应该运行
* 6. 循环迭代:updateCur1→AI active, updateCur2→AI skipped → AI 应该继续运行
* 7. 退出路径:reply 应该运行
*/
const nodes = [
createNode('start', FlowNodeTypeEnum.workflowStart),
createNode('ifElse1', FlowNodeTypeEnum.ifElseNode),
createNode('code', FlowNodeTypeEnum.code),
createNode('newArr', FlowNodeTypeEnum.variableUpdate),
createNode('getFirst', FlowNodeTypeEnum.code),
createNode('updateArr', FlowNodeTypeEnum.variableUpdate),
createNode('updateCur1', FlowNodeTypeEnum.variableUpdate),
createNode('AI', FlowNodeTypeEnum.chatNode),
createNode('updateCur2', FlowNodeTypeEnum.variableUpdate),
createNode('ifElse2', FlowNodeTypeEnum.ifElseNode),
createNode('updateHistory', FlowNodeTypeEnum.variableUpdate),
createNode('reply', FlowNodeTypeEnum.answerNode)
];
const edges = [
createEdge('start', 'ifElse1'),
// IF 分支(长链路)
createEdge('ifElse1', 'code', 'waiting', 'ifElse1-source-IF'),
createEdge('code', 'newArr'),
createEdge('newArr', 'getFirst'),
createEdge('getFirst', 'updateArr'),
createEdge('updateArr', 'updateCur1'),
createEdge('updateCur1', 'AI'),
// ELSE 分支(短链路)
createEdge('ifElse1', 'updateCur2', 'waiting', 'ifElse1-source-ELSE'),
createEdge('updateCur2', 'AI'),
// AI 后续
createEdge('AI', 'ifElse2'),
// 循环路径(IF 分支):更新历史后回到 getFirst
createEdge('ifElse2', 'updateHistory', 'waiting', 'ifElse2-source-IF'),
createEdge('updateHistory', 'getFirst'),
// 退出路径(ELSE 分支)
createEdge('ifElse2', 'reply', 'waiting', 'ifElse2-source-ELSE')
];
const edgeIndex = WorkflowQueue.buildEdgeIndex({ runtimeEdges: edges });
const edgeGroupsMap = WorkflowQueue.buildNodeEdgeGroupsMap({
runtimeNodes: nodes,
edgeIndex
});
it('getFirst 节点应该分成 2 组(初始 IF 路径 + 循环路径)', () => {
const groups = edgeGroupsMap.get('getFirst') || [];
expect(groups.length).toBe(2);
});
it('AI 节点应该分成 2 组(IF 分支 + ELSE 分支)', () => {
const groups = edgeGroupsMap.get('AI') || [];
expect(groups.length).toBe(2);
});
it('场景24.1: IF 分支首次执行,newArr→getFirst active,getFirst 应该运行', () => {
setEdgeStatus(edges, 'newArr', 'getFirst', 'active');
setEdgeStatus(edges, 'updateHistory', 'getFirst', 'waiting');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'getFirst')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(status).toBe('run');
});
it('场景24.2: IF 分支首次执行,updateCur1→AI active,updateCur2→AI skipped,AI 应该运行', () => {
setEdgeStatus(edges, 'updateCur1', 'AI', 'active');
setEdgeStatus(edges, 'updateCur2', 'AI', 'skipped');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'AI')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(status).toBe('run');
});
it('场景24.3: ELSE 分支执行,updateCur2→AI active,updateCur1→AI skipped,AI 应该运行', () => {
setEdgeStatus(edges, 'updateCur2', 'AI', 'active');
setEdgeStatus(edges, 'updateCur1', 'AI', 'skipped');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'AI')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(status).toBe('run');
});
it('场景24.4: ELSE 分支执行且退出,所有入边均 skipped,getFirst 应该跳过', () => {
// ELSE 路径:ifElse1 走 ELSE → updateCur2 → AI → ifElse2 走 ELSE → reply
// IF 链路(code/newArr/getFirst/updateArr/updateCur1)全部被 skipped
// ifElse2 走 ELSE 分支,updateHistory 也被 skipped,所以 updateHistory→getFirst 也为 skipped
setEdgeStatus(edges, 'newArr', 'getFirst', 'skipped');
setEdgeStatus(edges, 'updateHistory', 'getFirst', 'skipped');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'getFirst')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
// Group1[newArr→getFirst] 全部 skipped,Group2[updateHistory→getFirst] 全部 skipped
// getFirst 应该跳过
expect(status).toBe('skip');
});
it('场景24.5: 循环迭代,updateHistory→getFirst active,newArr→getFirst skipped,getFirst 应该运行', () => {
setEdgeStatus(edges, 'updateHistory', 'getFirst', 'active');
setEdgeStatus(edges, 'newArr', 'getFirst', 'skipped');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'getFirst')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(status).toBe('run');
});
it('场景24.6: 循环迭代,updateCur1→AI active,updateCur2→AI skipped,AI 应该继续运行', () => {
setEdgeStatus(edges, 'updateCur1', 'AI', 'active');
setEdgeStatus(edges, 'updateCur2', 'AI', 'skipped');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'AI')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(status).toBe('run');
});
it('场景24.7: 退出路径,ifElse2→reply active,reply 应该运行', () => {
setEdgeStatus(edges, 'ifElse2', 'reply', 'active');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'reply')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(status).toBe('run');
});
it('场景24.8: AI 两条边都 waiting,AI 应该等待', () => {
setEdgeStatus(edges, 'updateCur1', 'AI', 'waiting');
setEdgeStatus(edges, 'updateCur2', 'AI', 'waiting');
const status = WorkflowQueue.getNodeRunStatus({
node: nodes.find((n) => n.nodeId === 'AI')!,
nodeEdgeGroupsMap: edgeGroupsMap
});
expect(status).toBe('wait');
});
});
......@@ -31,3 +31,15 @@ export const createEdge = (
sourceHandle: sourceHandle || `${source}-source-right`,
targetHandle: targetHandle || `${target}-target-left`
});
export const setEdgeStatus = (
edges: RuntimeEdgeItemType[],
source: string,
target: string,
status: 'active' | 'waiting' | 'skipped'
) => {
const edge = edges.find((e) => e.source === source && e.target === target);
if (edge) {
edge.status = status;
}
};
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