Commit 22348594 by Ryo Committed by GitHub

perf: add process memory metrics (#6656)

* perf: reduce trace span and metrics

* perf: add process memory metrics

* fix: translations
parent 6e6b026d
import { configureMetricsFromEnv, disposeMetrics, getMeter } from '@fastgpt-sdk/otel/metrics'; import {
configureMetricsFromEnv,
disposeMetrics as disposeOtelMetrics,
getMeter
} from '@fastgpt-sdk/otel/metrics';
import { env } from '../../env'; import { env } from '../../env';
import { startRuntimeMetrics, stopRuntimeMetrics } from './runtime';
export async function configureMetrics() { export async function configureMetrics() {
await configureMetricsFromEnv({ await configureMetricsFromEnv({
...@@ -7,6 +12,13 @@ export async function configureMetrics() { ...@@ -7,6 +12,13 @@ export async function configureMetrics() {
defaultServiceName: 'fastgpt-client', defaultServiceName: 'fastgpt-client',
defaultMeterName: 'fastgpt-client' defaultMeterName: 'fastgpt-client'
}); });
startRuntimeMetrics();
}
export async function disposeMetrics() {
stopRuntimeMetrics();
await disposeOtelMetrics();
} }
export { disposeMetrics, getMeter }; export { getMeter };
import type {
BatchObservableCallback,
Meter,
Observable,
ObservableGauge
} from '@opentelemetry/api';
import { getMeter } from '@fastgpt-sdk/otel/metrics';
type RuntimeMetricAttributes = Record<string, never>;
type RuntimeObservableSet = {
meter: Meter;
processMemoryRss: ObservableGauge<RuntimeMetricAttributes>;
processMemoryHeapUsed: ObservableGauge<RuntimeMetricAttributes>;
processMemoryHeapTotal: ObservableGauge<RuntimeMetricAttributes>;
processMemoryExternal: ObservableGauge<RuntimeMetricAttributes>;
processMemoryArrayBuffers: ObservableGauge<RuntimeMetricAttributes>;
processUptime: ObservableGauge<RuntimeMetricAttributes>;
};
const prefix = 'fastgpt.runtime.process';
let runtimeMetricsRegistered = false;
let runtimeMeter: Meter | undefined;
let runtimeObservables: Observable<RuntimeMetricAttributes>[] = [];
let runtimeMetricsCallback: BatchObservableCallback<RuntimeMetricAttributes> | undefined;
function createRuntimeObservables(): RuntimeObservableSet {
const meter = getMeter('fastgpt.runtime');
return {
meter,
processMemoryRss: meter.createObservableGauge(`${prefix}.memory.rss`, {
description: 'Resident set size memory used by the current process',
unit: 'By'
}),
processMemoryHeapUsed: meter.createObservableGauge(`${prefix}.memory.heap_used`, {
description: 'V8 heap memory currently used by the current process',
unit: 'By'
}),
processMemoryHeapTotal: meter.createObservableGauge(`${prefix}.memory.heap_total`, {
description: 'Total V8 heap memory allocated for the current process',
unit: 'By'
}),
processMemoryExternal: meter.createObservableGauge(`${prefix}.memory.external`, {
description: 'Memory used by C++ objects bound to JavaScript objects',
unit: 'By'
}),
processMemoryArrayBuffers: meter.createObservableGauge(`${prefix}.memory.array_buffers`, {
description: 'Memory allocated for ArrayBuffer and SharedArrayBuffer instances',
unit: 'By'
}),
processUptime: meter.createObservableGauge(`${prefix}.uptime`, {
description: 'Process uptime',
unit: 's'
})
};
}
export function startRuntimeMetrics() {
if (runtimeMetricsRegistered) return;
const observables = createRuntimeObservables();
runtimeMeter = observables.meter;
runtimeObservables = [
observables.processMemoryRss,
observables.processMemoryHeapUsed,
observables.processMemoryHeapTotal,
observables.processMemoryExternal,
observables.processMemoryArrayBuffers,
observables.processUptime
];
runtimeMetricsCallback = (result) => {
const memoryUsage = process.memoryUsage();
result.observe(observables.processMemoryRss, memoryUsage.rss);
result.observe(observables.processMemoryHeapUsed, memoryUsage.heapUsed);
result.observe(observables.processMemoryHeapTotal, memoryUsage.heapTotal);
result.observe(observables.processMemoryExternal, memoryUsage.external);
result.observe(observables.processMemoryArrayBuffers, memoryUsage.arrayBuffers);
result.observe(observables.processUptime, process.uptime());
};
runtimeMeter.addBatchObservableCallback(runtimeMetricsCallback, runtimeObservables);
runtimeMetricsRegistered = true;
}
export function stopRuntimeMetrics() {
if (!runtimeMetricsRegistered || !runtimeMetricsCallback || !runtimeMeter) return;
runtimeMeter.removeBatchObservableCallback(runtimeMetricsCallback, runtimeObservables);
runtimeMetricsRegistered = false;
runtimeMeter = undefined;
runtimeObservables = [];
runtimeMetricsCallback = undefined;
}
...@@ -13,6 +13,37 @@ export type NextApiHandler<T = any> = ( ...@@ -13,6 +13,37 @@ export type NextApiHandler<T = any> = (
res: NextApiResponse<T> res: NextApiResponse<T>
) => unknown | Promise<unknown>; ) => unknown | Promise<unknown>;
function isIdLikeRouteSegment(segment: string) {
return (
/^\d{4,}$/.test(segment) ||
/^[0-9a-f]{24}$/i.test(segment) ||
/^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(segment) ||
/^[A-Za-z0-9_-]{16,}$/.test(segment)
);
}
function normalizeRouteSegment(segment: string) {
return isIdLikeRouteSegment(segment) ? ':id' : segment;
}
function parseHeaderNumber(value: string | string[] | undefined) {
const normalized = Array.isArray(value) ? value[0] : value;
if (!normalized) return undefined;
const parsed = Number(normalized);
return Number.isFinite(parsed) ? parsed : undefined;
}
function getRequestRoute(url: string) {
const [route = '/'] = url.split('?');
if (!route || route === '/') return '/';
return route
.split('/')
.map((segment) => normalizeRouteSegment(segment))
.join('/');
}
export const NextEntry = ({ export const NextEntry = ({
beforeCallback = [] beforeCallback = []
}: { }: {
...@@ -28,23 +59,22 @@ export const NextEntry = ({ ...@@ -28,23 +59,22 @@ export const NextEntry = ({
const responseLogger = getLogger(LogCategories.HTTP.RESPONSE); const responseLogger = getLogger(LogCategories.HTTP.RESPONSE);
const url = req.url || ''; const url = req.url || '';
const route = getRequestRoute(url);
const method = req.method?.toUpperCase() || ''; const method = req.method?.toUpperCase() || '';
const ip = req.headers['x-forwarded-for'] || req.socket?.remoteAddress; const ip = req.headers['x-forwarded-for'] || req.socket?.remoteAddress;
const userAgent = req.headers['user-agent']; const userAgent = req.headers['user-agent'];
const contentLength = req.headers['content-length']; const contentLength = req.headers['content-length'];
const requestBodySize = parseHeaderNumber(contentLength);
return withContext({ requestId }, async () => return withContext({ requestId }, async () =>
withActiveSpan( withActiveSpan(
{ {
name: `http.request ${method || 'UNKNOWN'} ${url || '/'}`, name: 'http.request',
tracerName: 'fastgpt.http', tracerName: 'fastgpt.http',
attributes: { attributes: {
'fastgpt.request.id': requestId,
'http.request.method': method, 'http.request.method': method,
'url.full': url, 'http.route': route,
'client.address': Array.isArray(ip) ? ip.join(',') : ip, 'http.request.body.size': requestBodySize
'user_agent.original': userAgent,
'http.request.body.size': contentLength
} }
}, },
async (span) => { async (span) => {
......
...@@ -29,6 +29,19 @@ export type ActiveSpanOptions = { ...@@ -29,6 +29,19 @@ export type ActiveSpanOptions = {
attributes?: Record<string, unknown>; attributes?: Record<string, unknown>;
}; };
const DEFAULT_PRODUCTION_TRACING_SAMPLE_RATIO = 0.05;
const DEFAULT_NON_PRODUCTION_TRACING_SAMPLE_RATIO = 1;
function getDefaultTracingSampleRatio() {
if (typeof env.TRACING_OTEL_SAMPLE_RATIO === 'number') {
return env.TRACING_OTEL_SAMPLE_RATIO;
}
return process.env.NODE_ENV === 'production'
? DEFAULT_PRODUCTION_TRACING_SAMPLE_RATIO
: DEFAULT_NON_PRODUCTION_TRACING_SAMPLE_RATIO;
}
function normalizeAttributes(attributes?: Record<string, unknown>) { function normalizeAttributes(attributes?: Record<string, unknown>) {
if (!attributes) return; if (!attributes) return;
...@@ -51,7 +64,7 @@ export async function configureTracing() { ...@@ -51,7 +64,7 @@ export async function configureTracing() {
env, env,
defaultServiceName: 'fastgpt-client', defaultServiceName: 'fastgpt-client',
defaultTracerName: 'fastgpt-client', defaultTracerName: 'fastgpt-client',
defaultSampleRatio: env.TRACING_OTEL_SAMPLE_RATIO defaultSampleRatio: getDefaultTracingSampleRatio()
}); });
} }
......
import { getNanoid } from '@fastgpt/global/common/string/tools'; import { getNanoid } from '@fastgpt/global/common/string/tools';
import { SpanStatusCode } from '@opentelemetry/api'; import { SpanStatusCode, trace, type Span } from '@opentelemetry/api';
import type { import type {
AIChatItemValueItemType, AIChatItemValueItemType,
ChatHistoryItemResType, ChatHistoryItemResType,
...@@ -63,13 +63,13 @@ import { TeamErrEnum } from '@fastgpt/global/common/error/code/team'; ...@@ -63,13 +63,13 @@ import { TeamErrEnum } from '@fastgpt/global/common/error/code/team';
import { i18nT } from '../../../../web/i18n/utils'; import { i18nT } from '../../../../web/i18n/utils';
import { validateFileUrlDomain } from '../../../common/security/fileUrlValidator'; import { validateFileUrlDomain } from '../../../common/security/fileUrlValidator';
import { classifyEdgesByDFS, findSCCs, isNodeInCycle, getEdgeType } from '../utils/tarjan'; import { classifyEdgesByDFS, findSCCs, isNodeInCycle, getEdgeType } from '../utils/tarjan';
import { observeWorkflowStep } from '../metrics'; import { observeWorkflowRun, observeWorkflowStep } from '../metrics';
import { withActiveSpan } from '../../../common/tracing'; import { withActiveSpan } from '../../../common/tracing';
const logger = getLogger(LogCategories.MODULE.WORKFLOW.DISPATCH);
import { delAgentRuntimeStopSign, shouldWorkflowStop } from './workflowStatus'; import { delAgentRuntimeStopSign, shouldWorkflowStop } from './workflowStatus';
import { runWithContext } from '../utils/context'; import { runWithContext } from '../utils/context';
const logger = getLogger(LogCategories.MODULE.WORKFLOW.DISPATCH);
type Props = Omit< type Props = Omit<
ChatDispatchProps, ChatDispatchProps,
'checkIsStopping' | 'workflowDispatchDeep' | 'timezone' | 'externalProvider' 'checkIsStopping' | 'workflowDispatchDeep' | 'timezone' | 'externalProvider'
...@@ -85,6 +85,68 @@ type NodeResponseCompleteType = Omit<NodeResponseType, 'responseData'> & { ...@@ -85,6 +85,68 @@ type NodeResponseCompleteType = Omit<NodeResponseType, 'responseData'> & {
[DispatchNodeResponseKeyEnum.nodeResponse]?: ChatHistoryItemResType; [DispatchNodeResponseKeyEnum.nodeResponse]?: ChatHistoryItemResType;
}; };
type WorkflowObservedStepResult = {
node: RuntimeNodeItemType;
runStatus: 'run';
result: NodeResponseCompleteType;
};
const tracedWorkflowStepTypes = new Set<FlowNodeTypeEnum>([
FlowNodeTypeEnum.appModule,
FlowNodeTypeEnum.pluginModule,
FlowNodeTypeEnum.agent,
FlowNodeTypeEnum.chatNode,
FlowNodeTypeEnum.datasetSearchNode,
FlowNodeTypeEnum.classifyQuestion,
FlowNodeTypeEnum.contentExtract,
FlowNodeTypeEnum.queryExtension,
FlowNodeTypeEnum.toolCall,
FlowNodeTypeEnum.httpRequest468,
FlowNodeTypeEnum.lafModule,
FlowNodeTypeEnum.code,
FlowNodeTypeEnum.readFiles,
FlowNodeTypeEnum.tool
]);
function shouldTraceWorkflowStep(nodeType: FlowNodeTypeEnum) {
return tracedWorkflowStepTypes.has(nodeType);
}
function getWorkflowStepStatus(result: WorkflowObservedStepResult): 'ok' | 'error' {
return result.result[DispatchNodeResponseKeyEnum.nodeResponse]?.error ? 'error' : 'ok';
}
function addWorkflowStepEvent({
eventName,
nodeType,
mode,
status,
durationMs
}: {
eventName: 'workflow.step.start' | 'workflow.step.end';
nodeType: FlowNodeTypeEnum;
mode: string;
status?: 'ok' | 'error';
durationMs?: number;
}) {
const activeSpan = trace.getActiveSpan();
if (!activeSpan) return;
const attributes: Record<string, string | number> = {
'fastgpt.workflow.node.type': nodeType,
'fastgpt.workflow.mode': mode
};
if (status) {
attributes['fastgpt.workflow.step.status'] = status;
}
if (typeof durationMs === 'number') {
attributes['fastgpt.workflow.step.duration_ms'] = durationMs;
}
activeSpan.addEvent(eventName, attributes);
}
// Run workflow // Run workflow
type WorkflowUsageProps = RequireOnlyOne<{ type WorkflowUsageProps = RequireOnlyOne<{
usageSource: UsageSourceEnum; usageSource: UsageSourceEnum;
...@@ -746,30 +808,13 @@ export class WorkflowQueue { ...@@ -746,30 +808,13 @@ export class WorkflowQueue {
runStatus: 'run'; runStatus: 'run';
result: NodeResponseCompleteType; result: NodeResponseCompleteType;
}> { }> {
const mode = this.isDebugMode ? 'test' : this.data.mode;
const stepMetricAttributes = { const stepMetricAttributes = {
workflowId: this.data.runningAppInfo.id,
workflowName: this.data.runningAppInfo.name,
nodeId: node.nodeId,
nodeName: node.name,
nodeType: node.flowNodeType, nodeType: node.flowNodeType,
mode: this.isDebugMode ? 'test' : this.data.mode mode
}; };
return observeWorkflowStep(stepMetricAttributes, () => const executeNode = async (stepSpan?: Span): Promise<WorkflowObservedStepResult> => {
withActiveSpan(
{
name: `workflow.step ${node.name || node.nodeId}`,
tracerName: 'fastgpt.workflow',
attributes: {
'fastgpt.workflow.id': this.data.runningAppInfo.id,
'fastgpt.workflow.name': this.data.runningAppInfo.name,
'fastgpt.workflow.node.id': node.nodeId,
'fastgpt.workflow.node.name': node.name,
'fastgpt.workflow.node.type': node.flowNodeType,
'fastgpt.workflow.mode': stepMetricAttributes.mode
}
},
async (stepSpan) => {
/* Inject data into module input */ /* Inject data into module input */
const getNodeRunParams = (node: RuntimeNodeItemType) => { const getNodeRunParams = (node: RuntimeNodeItemType) => {
if (node.flowNodeType === FlowNodeTypeEnum.pluginInput) { if (node.flowNodeType === FlowNodeTypeEnum.pluginInput) {
...@@ -856,7 +901,7 @@ export class WorkflowQueue { ...@@ -856,7 +901,7 @@ export class WorkflowQueue {
runtimeNodes: this.data.runtimeNodes, runtimeNodes: this.data.runtimeNodes,
runtimeEdges: this.data.runtimeEdges, runtimeEdges: this.data.runtimeEdges,
params, params,
mode: this.isDebugMode ? 'test' : this.data.mode mode
}; };
// run module // run module
...@@ -866,9 +911,7 @@ export class WorkflowQueue { ...@@ -866,9 +911,7 @@ export class WorkflowQueue {
const errorHandleId = getHandleId(node.nodeId, 'source_catch', 'right'); const errorHandleId = getHandleId(node.nodeId, 'source_catch', 'right');
try { try {
const result = (await callbackMap[node.flowNodeType]( const result = (await callbackMap[node.flowNodeType](dispatchData)) as NodeResponseType;
dispatchData
)) as NodeResponseType;
if (result.error) { if (result.error) {
// Run error and not catch error, skip all edges // Run error and not catch error, skip all edges
...@@ -891,18 +934,16 @@ export class WorkflowQueue { ...@@ -891,18 +934,16 @@ export class WorkflowQueue {
[DispatchNodeResponseKeyEnum.skipHandleId]: result[ [DispatchNodeResponseKeyEnum.skipHandleId]: result[
DispatchNodeResponseKeyEnum.skipHandleId DispatchNodeResponseKeyEnum.skipHandleId
] ]
? [ ? [...result[DispatchNodeResponseKeyEnum.skipHandleId], ...skipHandleIds].filter(
...result[DispatchNodeResponseKeyEnum.skipHandleId], Boolean
...skipHandleIds )
].filter(Boolean)
: skipHandleIds : skipHandleIds
}; };
} }
// Not error // Not error
const errorHandle = const errorHandle =
targetEdges.find((item) => item.sourceHandle === errorHandleId)?.sourceHandle || targetEdges.find((item) => item.sourceHandle === errorHandleId)?.sourceHandle || '';
'';
return { return {
...result, ...result,
...@@ -990,17 +1031,19 @@ export class WorkflowQueue { ...@@ -990,17 +1031,19 @@ export class WorkflowQueue {
// Error // Error
if (dispatchRes?.responseData?.error) { if (dispatchRes?.responseData?.error) {
if (stepSpan) {
stepSpan.setAttribute('fastgpt.workflow.step.error', true); stepSpan.setAttribute('fastgpt.workflow.step.error', true);
stepSpan.setStatus({ stepSpan.setStatus({
code: SpanStatusCode.ERROR, code: SpanStatusCode.ERROR,
message: String(dispatchRes.responseData.error) message: String(dispatchRes.responseData.error)
}); });
}
logger.warn('Workflow node returned error', { error: dispatchRes.responseData.error }); logger.warn('Workflow node returned error', { error: dispatchRes.responseData.error });
} else { } else if (stepSpan) {
stepSpan.setStatus({ code: SpanStatusCode.OK }); stepSpan.setStatus({ code: SpanStatusCode.OK });
} }
if (formatResponseData?.runningTime !== undefined) { if (stepSpan && formatResponseData?.runningTime !== undefined) {
stepSpan.setAttribute( stepSpan.setAttribute(
'fastgpt.workflow.step.running_time_seconds', 'fastgpt.workflow.step.running_time_seconds',
formatResponseData.runningTime formatResponseData.runningTime
...@@ -1015,8 +1058,65 @@ export class WorkflowQueue { ...@@ -1015,8 +1058,65 @@ export class WorkflowQueue {
[DispatchNodeResponseKeyEnum.nodeResponse]: formatResponseData [DispatchNodeResponseKeyEnum.nodeResponse]: formatResponseData
} }
}; };
};
if (shouldTraceWorkflowStep(node.flowNodeType)) {
return observeWorkflowStep(
stepMetricAttributes,
() =>
withActiveSpan(
{
name: 'workflow.step',
tracerName: 'fastgpt.workflow',
attributes: {
'fastgpt.workflow.node.type': node.flowNodeType,
'fastgpt.workflow.mode': mode
}
},
async (stepSpan) => executeNode(stepSpan)
),
{
getStatus: getWorkflowStepStatus
}
);
}
return observeWorkflowStep(
stepMetricAttributes,
async () => {
const stepStartedAt = Date.now();
addWorkflowStepEvent({
eventName: 'workflow.step.start',
nodeType: node.flowNodeType,
mode
});
try {
const result = await executeNode();
addWorkflowStepEvent({
eventName: 'workflow.step.end',
nodeType: node.flowNodeType,
mode,
status: getWorkflowStepStatus(result),
durationMs: Date.now() - stepStartedAt
});
return result;
} catch (error) {
addWorkflowStepEvent({
eventName: 'workflow.step.end',
nodeType: node.flowNodeType,
mode,
status: 'error',
durationMs: Date.now() - stepStartedAt
});
throw error;
}
},
{
getStatus: getWorkflowStepStatus
} }
)
); );
} }
private nodeRunWithSkip(node: RuntimeNodeItemType): { private nodeRunWithSkip(node: RuntimeNodeItemType): {
...@@ -1426,17 +1526,20 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR ...@@ -1426,17 +1526,20 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR
workflowId: data.runningAppInfo.id workflowId: data.runningAppInfo.id
}); });
return withActiveSpan( return observeWorkflowRun(
{
mode: data.mode,
isRoot: isRootRuntime
},
() =>
withActiveSpan(
{ {
name: isRootRuntime ? 'workflow.run' : 'workflow.child.run', name: isRootRuntime ? 'workflow.run' : 'workflow.child.run',
tracerName: 'fastgpt.workflow', tracerName: 'fastgpt.workflow',
attributes: { attributes: {
'fastgpt.workflow.id': data.runningAppInfo.id,
'fastgpt.workflow.name': data.runningAppInfo.name,
'fastgpt.workflow.mode': data.mode, 'fastgpt.workflow.mode': data.mode,
'fastgpt.workflow.depth': data.workflowDispatchDeep, 'fastgpt.workflow.depth': data.workflowDispatchDeep,
'fastgpt.workflow.is_root': isRootRuntime, 'fastgpt.workflow.is_root': isRootRuntime,
'fastgpt.workflow.chat_id': data.chatId,
'fastgpt.workflow.app_version': data.apiVersion, 'fastgpt.workflow.app_version': data.apiVersion,
'fastgpt.workflow.is_tool_call': !!data.isToolCall, 'fastgpt.workflow.is_tool_call': !!data.isToolCall,
'fastgpt.workflow.node_count': data.runtimeNodes.length, 'fastgpt.workflow.node_count': data.runtimeNodes.length,
...@@ -1543,6 +1646,10 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR ...@@ -1543,6 +1646,10 @@ export const runWorkflow = async (data: RunWorkflowProps): Promise<DispatchFlowR
durationSeconds durationSeconds
}; };
} }
),
{
getRunTimes: (result) => result[DispatchNodeResponseKeyEnum.runTimes]
}
); );
}; };
......
...@@ -2,30 +2,28 @@ import { getMeter } from '../../common/metrics'; ...@@ -2,30 +2,28 @@ import { getMeter } from '../../common/metrics';
type MetricAttributeValue = string | number | boolean; type MetricAttributeValue = string | number | boolean;
type MetricAttributes = Record<string, MetricAttributeValue>; type MetricAttributes = Record<string, MetricAttributeValue>;
type ObservationStatus = 'ok' | 'error';
export type WorkflowStepMetricAttributes = { type ObservationState = {
workflowId?: string; startedAt: bigint;
workflowName?: string; };
nodeId: string;
nodeName?: string; type ObserveMetricOptions<T> = {
nodeType: string; getStatus?: (result: T) => ObservationStatus;
};
export type WorkflowRunMetricAttributes = {
mode?: string; mode?: string;
isRoot?: boolean;
}; };
type ProcessSnapshot = { export type WorkflowStepMetricAttributes = {
rss: number; nodeType: string;
heapUsed: number; mode?: string;
external: number;
arrayBuffers: number;
cpuUser: number;
cpuSystem: number;
}; };
type StepObservationState = { type ObserveWorkflowRunOptions<T> = ObserveMetricOptions<T> & {
startedAt: bigint; getRunTimes?: (result: T) => number | undefined;
startSnapshot: ProcessSnapshot;
hadOverlapAtStart: boolean;
overlapVersionAtStart: number;
}; };
function normalizeAttributes(attributes: Record<string, unknown>): MetricAttributes { function normalizeAttributes(attributes: Record<string, unknown>): MetricAttributes {
...@@ -42,200 +40,129 @@ function normalizeAttributes(attributes: Record<string, unknown>): MetricAttribu ...@@ -42,200 +40,129 @@ function normalizeAttributes(attributes: Record<string, unknown>): MetricAttribu
return normalized; return normalized;
} }
function toMetricAttributes( function toRunMetricAttributes(
attributes: WorkflowRunMetricAttributes,
extras?: Record<string, unknown>
) {
return normalizeAttributes({
mode: attributes.mode,
is_root: attributes.isRoot,
...extras
});
}
function toStepMetricAttributes(
attributes: WorkflowStepMetricAttributes, attributes: WorkflowStepMetricAttributes,
extras?: Record<string, unknown> extras?: Record<string, unknown>
) { ) {
return normalizeAttributes({ return normalizeAttributes({
workflow_id: attributes.workflowId,
workflow_name: attributes.workflowName,
node_id: attributes.nodeId,
node_name: attributes.nodeName,
node_type: attributes.nodeType, node_type: attributes.nodeType,
mode: attributes.mode, mode: attributes.mode,
...extras ...extras
}); });
} }
function takeProcessSnapshot(): ProcessSnapshot { function beginObservation(): ObservationState {
const memory = process.memoryUsage();
const cpu = process.cpuUsage();
return { return {
rss: memory.rss, startedAt: process.hrtime.bigint()
heapUsed: memory.heapUsed,
external: memory.external,
arrayBuffers: memory.arrayBuffers,
cpuUser: cpu.user,
cpuSystem: cpu.system
}; };
} }
let activeWorkflowStepCount = 0; function getObservationDurationMs(state: ObservationState) {
let overlapVersion = 0; return Number(process.hrtime.bigint() - state.startedAt) / 1_000_000;
}
function beginStepObservation(): StepObservationState {
const state: StepObservationState = {
startedAt: process.hrtime.bigint(),
startSnapshot: takeProcessSnapshot(),
hadOverlapAtStart: activeWorkflowStepCount > 0,
overlapVersionAtStart: overlapVersion
};
activeWorkflowStepCount += 1; async function observeOperation<T>({
fn,
onStart,
onFinish,
options
}: {
fn: () => Promise<T> | T;
onStart?: () => void;
onFinish: (status: ObservationStatus, result: T | undefined, state: ObservationState) => void;
options?: ObserveMetricOptions<T>;
}): Promise<T> {
const observationState = beginObservation();
onStart?.();
if (activeWorkflowStepCount > 1) { try {
overlapVersion += 1; const result = await fn();
const status = options?.getStatus?.(result) ?? 'ok';
onFinish(status, result, observationState);
return result;
} catch (error) {
onFinish('error', undefined, observationState);
throw error;
} }
return state;
} }
const meter = getMeter('fastgpt.workflow'); const meter = getMeter('fastgpt.workflow');
const prefix = 'fastgpt.workflow'; const prefix = 'fastgpt.workflow';
const stepDuration = meter.createHistogram(`${prefix}.step.duration`, { const runDuration = meter.createHistogram(`${prefix}.run.duration`, {
description: 'Workflow step execution duration', description: 'Workflow run duration',
unit: 'ms' unit: 'ms'
}); });
const stepExecutions = meter.createCounter(`${prefix}.step.executions`, { const runExecutions = meter.createCounter(`${prefix}.run.count`, {
description: 'Workflow step execution count' description: 'Workflow run 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`, { const runActive = meter.createUpDownCounter(`${prefix}.run.active`, {
description: 'Workflow process array buffer memory snapshot at step end', description: 'Workflow runs currently executing'
unit: 'By'
}); });
const stepMemoryRssGrowth = meter.createHistogram(`${prefix}.step.memory.rss_growth`, { const runTimes = meter.createHistogram(`${prefix}.run.run_times`, {
description: 'Workflow process RSS memory growth during non-overlapping step execution', description: 'Workflow total run times before completion'
unit: 'By'
}); });
const stepMemoryHeapUsedGrowth = meter.createHistogram(`${prefix}.step.memory.heap_used_growth`, { const stepDuration = meter.createHistogram(`${prefix}.step.duration`, {
description: 'Workflow process heap used memory growth during non-overlapping step execution', description: 'Workflow step execution duration',
unit: 'By' unit: 'ms'
}); });
const stepMemoryExternalGrowth = meter.createHistogram(`${prefix}.step.memory.external_growth`, { const stepExecutions = meter.createCounter(`${prefix}.step.count`, {
description: 'Workflow process external memory growth during non-overlapping step execution', description: 'Workflow step execution count'
unit: 'By'
}); });
export async function observeWorkflowStep<T>( export async function observeWorkflowRun<T>(
attributes: WorkflowStepMetricAttributes, attributes: WorkflowRunMetricAttributes,
fn: () => Promise<T> | T fn: () => Promise<T> | T,
options?: ObserveWorkflowRunOptions<T>
): Promise<T> { ): Promise<T> {
const observationState = beginStepObservation(); const baseAttributes = toRunMetricAttributes(attributes);
const baseAttributes = toMetricAttributes(attributes);
return observeOperation({
stepActive.add(1, baseAttributes); fn,
options,
onStart: () => {
runActive.add(1, baseAttributes);
},
onFinish: (status, result, state) => {
const metricAttributes = toRunMetricAttributes(attributes, { status });
runDuration.record(getObservationDurationMs(state), metricAttributes);
runExecutions.add(1, metricAttributes);
const workflowRunTimes = result ? options?.getRunTimes?.(result) : undefined;
if (typeof workflowRunTimes === 'number' && Number.isFinite(workflowRunTimes)) {
runTimes.record(workflowRunTimes, metricAttributes);
}
try { runActive.add(-1, baseAttributes);
const result = await fn();
recordWorkflowStepEnd(attributes, observationState, 'ok', baseAttributes);
return result;
} catch (error) {
recordWorkflowStepEnd(attributes, observationState, 'error', baseAttributes);
throw error;
} }
});
} }
function recordWorkflowStepEnd( export async function observeWorkflowStep<T>(
attributes: WorkflowStepMetricAttributes, attributes: WorkflowStepMetricAttributes,
observationState: StepObservationState, fn: () => Promise<T> | T,
status: 'ok' | 'error', options?: ObserveMetricOptions<T>
baseAttributes: MetricAttributes ): Promise<T> {
) { return observeOperation({
const endSnapshot = takeProcessSnapshot(); fn,
const metricAttributes = toMetricAttributes(attributes, { status }); options,
const stepOverlap = onFinish: (status, _result, state) => {
observationState.hadOverlapAtStart || observationState.overlapVersionAtStart !== overlapVersion; const metricAttributes = toStepMetricAttributes(attributes, { status });
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); stepDuration.record(getObservationDurationMs(state), metricAttributes);
stepExecutions.add(1, 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);
} }
...@@ -16,6 +16,14 @@ type FileTypeSelectorValue = { ...@@ -16,6 +16,14 @@ type FileTypeSelectorValue = {
customFileExtensionList?: string[]; customFileExtensionList?: string[];
}; };
const fileExtensionTypeTranslationMap = new Map<FileExtensionKeyType, string>([
['canSelectFile', 'app:upload_file_extension_type_canSelectFile'],
['canSelectImg', 'app:upload_file_extension_type_canSelectImg'],
['canSelectVideo', 'app:upload_file_extension_type_canSelectVideo'],
['canSelectAudio', 'app:upload_file_extension_type_canSelectAudio'],
['canSelectCustomFileExtension', 'app:upload_file_extension_type_canSelectCustomFileExtension']
]);
export const FileTypeSelectorPanel = ({ export const FileTypeSelectorPanel = ({
value, value,
onChange onChange
...@@ -190,7 +198,7 @@ export const FileTypeSelectorPanel = ({ ...@@ -190,7 +198,7 @@ export const FileTypeSelectorPanel = ({
onChange={(e) => handleTypeChange(type as FileExtensionKeyType, e.target.checked)} onChange={(e) => handleTypeChange(type as FileExtensionKeyType, e.target.checked)}
> >
<Box color={'myGray.900'} lineHeight={1}> <Box color={'myGray.900'} lineHeight={1}>
{t(`app:upload_file_extension_type_${type}`)} {t(fileExtensionTypeTranslationMap.get(type as FileExtensionKeyType) || type)}
</Box> </Box>
<Box mt={1} fontSize={'xs'} color={'myGray.500'} wordBreak={'break-word'} w="full"> <Box mt={1} fontSize={'xs'} color={'myGray.500'} wordBreak={'break-word'} w="full">
{exts.map((ext) => ext.slice(1)).join('/')} {exts.map((ext) => ext.slice(1)).join('/')}
......
...@@ -466,6 +466,10 @@ ...@@ -466,6 +466,10 @@
"upload_file_extension_type_canSelectCustomFileExtension": "Custom file extension type", "upload_file_extension_type_canSelectCustomFileExtension": "Custom file extension type",
"upload_file_extension_type_canSelectCustomFileExtension_placeholder": "file extension name", "upload_file_extension_type_canSelectCustomFileExtension_placeholder": "file extension name",
"upload_file_extension_types": "Supported file types", "upload_file_extension_types": "Supported file types",
"upload_file_extension_type_canSelectAudio": "Audio",
"upload_file_extension_type_canSelectFile": "Document",
"upload_file_extension_type_canSelectImg": "Image",
"upload_file_extension_type_canSelectVideo": "Video",
"upload_file_max_amount": "Maximum File Quantity", "upload_file_max_amount": "Maximum File Quantity",
"upload_file_max_amount_tip": "Maximum number of files uploaded in a single round of conversation", "upload_file_max_amount_tip": "Maximum number of files uploaded in a single round of conversation",
"upload_method": "Upload method", "upload_method": "Upload method",
......
...@@ -468,6 +468,10 @@ ...@@ -468,6 +468,10 @@
"upload_file_extension_types": "支持上传的类型", "upload_file_extension_types": "支持上传的类型",
"upload_file_max_amount": "最大文件数量", "upload_file_max_amount": "最大文件数量",
"upload_file_max_amount_tip": "单轮对话中最大上传文件数量", "upload_file_max_amount_tip": "单轮对话中最大上传文件数量",
"upload_file_extension_type_canSelectAudio": "音频",
"upload_file_extension_type_canSelectFile": "文档",
"upload_file_extension_type_canSelectImg": "图片",
"upload_file_extension_type_canSelectVideo": "视频",
"upload_method": "上传方式", "upload_method": "上传方式",
"url_upload": "文件链接", "url_upload": "文件链接",
"use_agent_sandbox": "虚拟机", "use_agent_sandbox": "虚拟机",
......
...@@ -452,6 +452,10 @@ ...@@ -452,6 +452,10 @@
"upload_file_extension_type_canSelectCustomFileExtension": "自定義文件擴展類型", "upload_file_extension_type_canSelectCustomFileExtension": "自定義文件擴展類型",
"upload_file_extension_type_canSelectCustomFileExtension_placeholder": "文件擴展名", "upload_file_extension_type_canSelectCustomFileExtension_placeholder": "文件擴展名",
"upload_file_extension_types": "支持上傳的類型", "upload_file_extension_types": "支持上傳的類型",
"upload_file_extension_type_canSelectAudio": "音頻",
"upload_file_extension_type_canSelectFile": "文檔",
"upload_file_extension_type_canSelectImg": "圖片",
"upload_file_extension_type_canSelectVideo": "視頻",
"upload_file_max_amount": "最大檔案數量", "upload_file_max_amount": "最大檔案數量",
"upload_file_max_amount_tip": "單輪對話中最大上傳檔案數量", "upload_file_max_amount_tip": "單輪對話中最大上傳檔案數量",
"upload_method": "上傳方式", "upload_method": "上傳方式",
......
...@@ -6,6 +6,8 @@ import withRspack from 'next-rspack'; ...@@ -6,6 +6,8 @@ import withRspack from 'next-rspack';
const withBundleAnalyzer = withBundleAnalyzerInit({ enabled: process.env.ANALYZE === 'true' }); const withBundleAnalyzer = withBundleAnalyzerInit({ enabled: process.env.ANALYZE === 'true' });
const isDev = process.env.NODE_ENV === 'development'; const isDev = process.env.NODE_ENV === 'development';
const isWebpack = process.env.WEBPACK === '1';
const isRspack = isDev && !isWebpack;
const nextConfig: NextConfig = { const nextConfig: NextConfig = {
basePath: process.env.NEXT_PUBLIC_BASE_URL, basePath: process.env.NEXT_PUBLIC_BASE_URL,
...@@ -218,8 +220,5 @@ const nextConfig: NextConfig = { ...@@ -218,8 +220,5 @@ const nextConfig: NextConfig = {
} }
}; };
const configWithPluginsExceptWithRspack = withBundleAnalyzer(nextConfig); const config = withBundleAnalyzer(nextConfig);
export default isRspack ? withRspack(config) : config;
export default isDev
? withRspack(configWithPluginsExceptWithRspack)
: configWithPluginsExceptWithRspack;
...@@ -4,6 +4,7 @@ ...@@ -4,6 +4,7 @@
"private": false, "private": false,
"scripts": { "scripts": {
"dev": "NODE_OPTIONS='--max-old-space-size=8192' npm run build:workers && next dev", "dev": "NODE_OPTIONS='--max-old-space-size=8192' npm run build:workers && next dev",
"dev:webpack": "NODE_OPTIONS='--max-old-space-size=8192' npm run build:workers && WEBPACK=1 next dev --webpack",
"build": "npm run build:workers && next build --debug --webpack", "build": "npm run build:workers && next build --debug --webpack",
"start": "next start", "start": "next start",
"build:workers": "npx tsx scripts/build-workers.ts", "build:workers": "npx tsx scripts/build-workers.ts",
......
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