Commit 778fb386 by Archer Committed by GitHub

feat: add liteparse (#7069)

* feat: add liteparse

* perf code

* perf code

* perf code

* fix: load liteparse as esm in worker

* fix: add liteparse runtime libs to app image

* perf: reduce read file worker memory usage

* fix: use compatible worker transfer type

* fix: respect read file worker concurrency
parent f2a89141
......@@ -47,6 +47,7 @@ fastgpt-plugin:
4. 输入引导配置增加校验,避免错误配置了自定义词库地址。
5. 工作流数组引用类型增强校验,避免刚好与二维数据冲突。
6. 知识库被删除后,应用编排时优雅提示。
7. PDF 解析,将 PDFJs 替换成 `liteparse`,速度提高 3 倍。
## 🐛 修复
......
......@@ -28,7 +28,6 @@ export const readRawTextByLocalFile = async (params: readRawTextByLocalFileParam
const { path } = params;
const extension = path?.split('.')?.pop()?.toLowerCase() || '';
const buffer = await fs.promises.readFile(path);
return readFileContentByBuffer({
......@@ -97,7 +96,7 @@ export const readFileContentByBuffer = async ({
const { data: response } = await axios.post<{
pages: number;
markdown: string;
error?: Object | string;
error?: object | string;
}>(url, data, {
timeout: 600000,
headers: {
......@@ -188,12 +187,14 @@ export const readFileContentByBuffer = async ({
const start = Date.now();
logger.debug('Start parsing file', { extension });
let { rawText, formatText, imageList } = await (async () => {
const parsedFile = await (async () => {
if (extension === 'pdf') {
return await pdfParseFn();
}
return await systemParse();
})();
const { imageList } = parsedFile;
let { rawText, formatText } = parsedFile;
logger.debug('File parsing completed', { extension, durationMs: Date.now() - start });
......
......@@ -19,6 +19,7 @@
"@fastgpt-sdk/sandbox-adapter": "workspace:*",
"@fastgpt-sdk/storage": "workspace:*",
"@fastgpt/global": "workspace:*",
"@llamaindex/liteparse": "catalog:",
"@mariozechner/pi-agent-core": "^0.67.3",
"@mariozechner/pi-ai": "^0.67.3",
"@maxmind/geoip2-node": "^6.3.4",
......@@ -66,7 +67,6 @@
"node-xlsx": "^0.24.0",
"p-limit": "^7.2.0",
"papaparse": "5.4.1",
"pdfjs-dist": "4.10.38",
"pg": "^8.10.0",
"pino": "^9.7.0",
"pino-opentelemetry-transport": "^1.0.1",
......
......@@ -101,6 +101,11 @@ describe('readRawTextByLocalFile', () => {
});
expect(result.rawText).toBe('Hello World');
expect(mockReadRawContentFromBuffer).toHaveBeenLastCalledWith(
expect.objectContaining({
extension: 'txt'
})
);
});
it('should extract extension from file path', async () => {
......@@ -135,6 +140,11 @@ describe('readFileContentByBuffer', () => {
});
expect(result.rawText).toBe('Hello from buffer');
expect(mockReadRawContentFromBuffer).toHaveBeenLastCalledWith(
expect.objectContaining({
extension: 'txt'
})
);
});
it('should use system parse for non-pdf files', async () => {
......
......@@ -80,8 +80,10 @@ describe('worker/function', () => {
mockEnv.PARSE_FILE_TIMEOUT_SECONDS = 300;
});
it('使用 SharedArrayBuffer 包装 Buffer 并通过 pool.run 派发', async () => {
const original = Buffer.from('hello world', 'utf-8');
it('默认 transfer 独占 Buffer 并通过 pool.run 派发', async () => {
const original = Buffer.allocUnsafeSlow(11);
original.write('hello world', 'utf-8');
const sourceArrayBuffer = original.buffer;
const expected = { rawText: 'parsed-content' };
mockRun.mockResolvedValueOnce(expected);
......@@ -107,12 +109,28 @@ describe('worker/function', () => {
expect(runArg.extension).toBe('txt');
expect(runArg.encoding).toBe('utf-8');
expect(runArg.bufferSize).toBe(original.length);
expect(runArg.sharedBuffer).toBeInstanceOf(SharedArrayBuffer);
expect(runArg.buffer).toBe(sourceArrayBuffer);
expect(runArg.sharedBuffer).toBeUndefined();
expect(mockRun.mock.calls[0][1]).toEqual([sourceArrayBuffer]);
});
it('Buffer 不独占 ArrayBuffer 时回退到 SharedArrayBuffer', async () => {
const original = Buffer.from('prefix:hello world').subarray('prefix:'.length);
mockRun.mockResolvedValueOnce({ rawText: 'parsed-content' });
await readRawContentFromBuffer({
extension: 'txt',
encoding: 'utf-8',
buffer: original
});
// SharedArrayBuffer 内容必须完整复刻原 Buffer
expect(runArg.sharedBuffer.byteLength).toBe(original.length);
const sharedView = new Uint8Array(runArg.sharedBuffer);
expect(Array.from(sharedView)).toEqual(Array.from(original));
const runArg = mockRun.mock.calls[0][0];
expect(runArg.buffer).toBeUndefined();
expect(runArg.sharedBuffer).toBeInstanceOf(SharedArrayBuffer);
expect(mockRun.mock.calls[0][1]).toBeUndefined();
expect(Buffer.from(new Uint8Array(runArg.sharedBuffer)).toString('utf-8')).toBe(
'hello world'
);
});
it('空 Buffer 也能正常构造(byteLength 为 0)', async () => {
......@@ -127,7 +145,9 @@ describe('worker/function', () => {
expect(result).toEqual({ rawText: '' });
const runArg = mockRun.mock.calls[0][0];
expect(runArg.bufferSize).toBe(0);
expect(runArg.sharedBuffer.byteLength).toBe(0);
expect(runArg.buffer.byteLength).toBe(0);
expect(runArg.sharedBuffer).toBeUndefined();
expect(mockRun.mock.calls[0][1]).toEqual([runArg.buffer]);
});
it('二进制 Buffer 不应在拷贝过程中失真', async () => {
......@@ -142,8 +162,8 @@ describe('worker/function', () => {
});
const runArg = mockRun.mock.calls[0][0];
const sharedView = new Uint8Array(runArg.sharedBuffer);
expect(Array.from(sharedView)).toEqual(Array.from(bytes));
const view = new Uint8Array(runArg.buffer ?? runArg.sharedBuffer);
expect(Array.from(view)).toEqual(Array.from(bytes));
});
it('PARSE_FILE_WORKERS 自定义值生效', async () => {
......@@ -186,6 +206,55 @@ describe('worker/function', () => {
).rejects.toThrow('parse failed');
});
it('并发文件解析直接交给 readFile worker pool,并发数由 PARSE_FILE_WORKERS 决定', async () => {
let activeCount = 0;
let maxActiveCount = 0;
const callOrder: string[] = [];
mockEnv.PARSE_FILE_WORKERS = 3;
mockRun.mockImplementation(
async (props: {
extension: string;
buffer?: ArrayBuffer;
sharedBuffer?: SharedArrayBuffer;
}) => {
activeCount += 1;
maxActiveCount = Math.max(maxActiveCount, activeCount);
const rawBuffer = props.buffer ?? props.sharedBuffer;
expect(rawBuffer).toBeDefined();
callOrder.push(Buffer.from(new Uint8Array(rawBuffer!)).toString('utf-8'));
await new Promise((resolve) => setTimeout(resolve, 20));
activeCount -= 1;
return { rawText: 'ok' };
}
);
const results = await Promise.all([
readRawContentFromBuffer({
extension: 'pdf',
encoding: 'utf-8',
buffer: Buffer.from('pdf-1')
}),
readRawContentFromBuffer({
extension: 'txt',
encoding: 'utf-8',
buffer: Buffer.from('txt-1')
}),
readRawContentFromBuffer({
extension: 'md',
encoding: 'utf-8',
buffer: Buffer.from('md-1')
})
]);
expect(results).toEqual([{ rawText: 'ok' }, { rawText: 'ok' }, { rawText: 'ok' }]);
expect(mockRun).toHaveBeenCalledTimes(3);
expect(maxActiveCount).toBe(3);
expect(callOrder).toEqual(expect.arrayContaining(['pdf-1', 'txt-1', 'md-1']));
});
it('多次调用每次都通过 getWorkerController 获取池(不在本层缓存)', async () => {
mockRun.mockResolvedValue({ rawText: 'ok' });
......@@ -210,18 +279,18 @@ describe('worker/function', () => {
expect(mockRun).toHaveBeenCalledTimes(3);
});
it('每次调用都生成新的 SharedArrayBuffer(避免跨任务串扰)', async () => {
it('fallback 路径每次调用都生成新的 SharedArrayBuffer(避免跨任务串扰)', async () => {
mockRun.mockResolvedValue({ rawText: 'ok' });
await readRawContentFromBuffer({
extension: 'txt',
encoding: 'utf-8',
buffer: Buffer.from('aaa')
buffer: Buffer.from('xaaa').subarray(1)
});
await readRawContentFromBuffer({
extension: 'txt',
encoding: 'utf-8',
buffer: Buffer.from('bbb')
buffer: Buffer.from('xbbb').subarray(1)
});
const sab1 = mockRun.mock.calls[0][0].sharedBuffer;
......
import { beforeEach, describe, expect, it, vi } from 'vitest';
const { mockLiteParseParse, mockLiteParseConstructor } = vi.hoisted(() => ({
mockLiteParseParse: vi.fn(),
mockLiteParseConstructor: vi.fn()
}));
vi.mock('@llamaindex/liteparse', () => ({
LiteParse: vi.fn(function MockLiteParse(config) {
mockLiteParseConstructor(config);
return {
parse: mockLiteParseParse
};
})
}));
import { readPdfFile } from '@fastgpt/service/worker/readFile/extension/pdf';
describe('readPdfFile', () => {
beforeEach(() => {
vi.clearAllMocks();
});
it('分批使用 LiteParse,并对全部 textItems 做统一后处理', async () => {
mockLiteParseParse
.mockResolvedValueOnce({
text: 'fallback text 1',
pages: [
{
pageNum: 1,
width: 1000,
height: 1000,
text: '',
textItems: [
{
text: '人工智能正在快速发展并进入规模化落地阶段',
x: 80,
y: 100,
width: 240,
height: 12
}
]
}
]
})
.mockResolvedValueOnce({
text: 'fallback text 2',
pages: [
{
pageNum: 101,
width: 1000,
height: 1000,
text: '',
textItems: [{ text: '并推动产业升级。', x: 80, y: 120, width: 96, height: 12 }]
}
]
})
.mockResolvedValueOnce({
text: '',
pages: []
});
const result = await readPdfFile({
extension: 'pdf',
encoding: 'utf-8',
buffer: Buffer.from('pdf')
});
expect(result.rawText).toBe('人工智能正在快速发展并进入规模化落地阶段并推动产业升级。\n');
expect(mockLiteParseParse).toHaveBeenCalledTimes(3);
expect(mockLiteParseConstructor.mock.calls.map(([config]) => config.targetPages)).toEqual([
'1-100',
'101-200',
'201-300'
]);
});
it('LiteParse 失败时直接抛出错误', async () => {
const error = new Error('native load failed');
mockLiteParseParse.mockRejectedValue(error);
await expect(
readPdfFile({
extension: 'pdf',
encoding: 'utf-8',
buffer: Buffer.from('pdf')
})
).rejects.toThrow(error);
});
});
import { describe, expect, it } from 'vitest';
import {
extractPageLines,
postprocessLiteParsePages
} from '@fastgpt/service/worker/readFile/utils/pdfTextPostprocess';
const textItem = ({
text,
x = 80,
y,
width,
height = 12,
fontSize = 12
}: {
text: string;
x?: number;
y: number;
width?: number;
height?: number;
fontSize?: number;
}) => ({
text,
x,
y,
width: width ?? text.length * 12,
height,
fontSize
});
describe('pdfTextPostprocess', () => {
it('按坐标重组同一行,并保守合并中文视觉换行', () => {
const text = postprocessLiteParsePages([
{
height: 1000,
textItems: [
textItem({ text: 'AI', x: 80, y: 100, width: 14 }),
textItem({ text: '技术正在快速发展,带动产业链上下游形成新的增长空间', x: 102, y: 100 }),
textItem({ text: '也对数据治理、算力供给和模型安全提出更高要求。', y: 120 })
]
}
]);
expect(text).toBe(
'AI 技术正在快速发展,带动产业链上下游形成新的增长空间也对数据治理、算力供给和模型安全提出更高要求。\n'
);
});
it('保留标题、列表和目录行的段落边界', () => {
const text = postprocessLiteParsePages([
{
height: 1000,
textItems: [
textItem({ text: '1.1 发展背景', y: 100 }),
textItem({ text: '人工智能产业已经进入规模化落地阶段。', y: 120 }),
textItem({ text: '(一)算力基础设施', y: 160 }),
textItem({ text: '目录章节................ 12', y: 200 })
]
}
]);
expect(text).toBe(
'1.1 发展背景\n\n人工智能产业已经进入规模化落地阶段。\n\n(一)算力基础设施\n\n目录章节................ 12\n'
);
});
it('过滤页眉页脚和纯页码', () => {
const lines = extractPageLines({
height: 1000,
textItems: [
textItem({ text: '页眉噪声', y: 20 }),
textItem({ text: '正文内容。', y: 120 }),
textItem({ text: '42', y: 930 }),
textItem({ text: '页脚噪声', y: 980 })
]
});
const text = postprocessLiteParsePages([
{
height: 1000,
textItems: lines.map((text, index) => textItem({ text, y: 100 + index * 20 }))
}
]);
expect(lines).toEqual(['正文内容。', '42']);
expect(text).toBe('正文内容。\n');
});
it('只删除重复噪声整行,不把普通短词从正文中删除', () => {
const text = postprocessLiteParsePages([
{
height: 1000,
textItems: [
textItem({ text: '操作', y: 100 }),
textItem({ text: '操作步骤如下,用户可以按需配置。', y: 120 })
]
},
{
height: 1000,
textItems: [textItem({ text: '操作', y: 100 }), textItem({ text: '第二页正文。', y: 120 })]
},
{
height: 1000,
textItems: [textItem({ text: '操作', y: 100 }), textItem({ text: '第三页正文。', y: 120 })]
}
]);
expect(text).toContain('操作步骤如下,用户可以按需配置。');
expect(text).not.toContain('\n\n操作\n\n');
});
});
import { describe, it, expect, beforeAll, afterEach, afterAll, vi } from 'vitest';
import path from 'path';
import { existsSync } from 'fs';
import { existsSync, readFileSync } from 'fs';
/*
* 真实 spawn 测试:使用 projects/app/worker/readFile.js 构建产物,
......@@ -17,6 +17,12 @@ const REAL_WORKER_PATH = path.join(APP_PROJECT_DIR, 'worker/readFile.js');
const shouldRunIntegration =
process.env.RUN_READ_FILE_WORKER_INTEGRATION === 'true' && existsSync(REAL_WORKER_PATH);
const pdfFixturePath = process.env.RUN_READ_FILE_WORKER_PDF_PATH;
const itIfPdfFixture = pdfFixturePath && existsSync(pdfFixturePath) ? it : it.skip;
const shouldRunPdfStress =
process.env.RUN_READ_FILE_WORKER_PDF_STRESS === 'true' &&
Boolean(pdfFixturePath && existsSync(pdfFixturePath));
const itIfPdfStress = shouldRunPdfStress ? it : it.skip;
const { WorkerNameEnum } = await import('@fastgpt/service/worker/utils');
const { readRawContentFromBuffer } = await import('@fastgpt/service/worker/function');
......@@ -43,6 +49,11 @@ const parseText = (text: string) =>
buffer: Buffer.from(text, 'utf-8')
});
const getPositiveIntegerEnv = (name: string, defaultValue: number) => {
const value = Number(process.env[name]);
return Number.isInteger(value) && value > 0 ? value : defaultValue;
};
const destroyReadFilePool = async () => {
const workerPoll = (global as any).workerPoll;
const pool = workerPoll?.[WorkerNameEnum.readFile];
......@@ -63,7 +74,6 @@ describeIfEnabled('readFile worker (real spawn integration)', () => {
let cwdSpy: ReturnType<typeof vi.spyOn>;
if (process.env.RUN_READ_FILE_WORKER_INTEGRATION === 'true' && !existsSync(REAL_WORKER_PATH)) {
// eslint-disable-next-line no-console
console.warn(
`[skipped] readFile worker integration requires RUN_READ_FILE_WORKER_INTEGRATION=true and worker bundle at ${REAL_WORKER_PATH}.`
);
......@@ -113,6 +123,89 @@ describeIfEnabled('readFile worker (real spawn integration)', () => {
expect(result.rawText).toContain('Shanghai');
});
itIfPdfFixture(
'解析 pdf(真实 worker + LiteParse)',
async () => {
const result = await readRawContentFromBuffer({
extension: 'pdf',
encoding: 'utf-8',
buffer: readFileSync(pdfFixturePath!)
});
expect(result.rawText.length).toBeGreaterThan(1000);
expect(result.rawText).toContain('人工智能');
},
60000
);
itIfPdfFixture(
'并发 pdf 直接交给真实 worker pool,按 PARSE_FILE_WORKERS 控制并发',
async () => {
const concurrency = 4;
const fileBuffer = readFileSync(pdfFixturePath!);
const results = await Promise.all(
Array.from({ length: concurrency }, () =>
readRawContentFromBuffer({
extension: 'pdf',
encoding: 'utf-8',
buffer: Buffer.from(fileBuffer)
})
)
);
expect(results).toHaveLength(concurrency);
results.forEach((result) => {
expect(result.rawText.length).toBeGreaterThan(1000);
expect(result.rawText).toContain('人工智能');
});
},
120000
);
itIfPdfStress(
'pdf worker 压测:多轮并发提交给 worker pool 后稳定返回',
async () => {
const concurrency = getPositiveIntegerEnv('RUN_READ_FILE_WORKER_PDF_STRESS_CONCURRENCY', 4);
const rounds = getPositiveIntegerEnv('RUN_READ_FILE_WORKER_PDF_STRESS_ROUNDS', 5);
const fileBuffer = readFileSync(pdfFixturePath!);
const durations: number[] = [];
const startedAt = Date.now();
for (let round = 0; round < rounds; round++) {
const roundStartedAt = Date.now();
const results = await Promise.all(
Array.from({ length: concurrency }, () =>
readRawContentFromBuffer({
extension: 'pdf',
encoding: 'utf-8',
buffer: Buffer.from(fileBuffer)
})
)
);
durations.push(Date.now() - roundStartedAt);
results.forEach((result) => {
expect(result.rawText.length).toBeGreaterThan(1000);
expect(result.rawText).toContain('人工智能');
});
}
const pool = getReadFilePool();
expect(pool.workerQueue.length).toBeLessThanOrEqual(pool.maxReservedThreads);
console.info('pdf worker stress summary', {
concurrency,
rounds,
totalTasks: concurrency * rounds,
wallMs: Date.now() - startedAt,
roundMs: durations,
workerCount: pool.workerQueue.length
});
},
120000
);
it('未知扩展名应被 reject', async () => {
await expect(
readRawContentFromBuffer({
......@@ -146,7 +239,7 @@ describeIfEnabled('readFile worker (real spawn integration)', () => {
expect(sameWorker?.tasksCompleted).toBe(initialTasks + 5);
});
it('并发场景:池按上限扩容,所有任务都成功返回', async () => {
it('并发场景:readFile 入口直接交给 worker pool,所有任务都成功返回', async () => {
const concurrency = 4;
const results = await Promise.all(
......@@ -163,7 +256,6 @@ describeIfEnabled('readFile worker (real spawn integration)', () => {
results.forEach((r, i) => expect(r.rawText).toBe(`payload-${i}`));
const pool = getReadFilePool();
// 池子大小不应超过 maxReservedThreads
expect(pool.workerQueue.length).toBeLessThanOrEqual(pool.maxReservedThreads);
expect(pool.workerQueue.length).toBeGreaterThan(1);
});
......
......@@ -24,7 +24,8 @@ export const text2Chunks = (props: SplitProps) => {
type ReadFileWorkerProps = {
extension: string;
encoding: string;
sharedBuffer: SharedArrayBuffer;
buffer?: ArrayBuffer;
sharedBuffer?: SharedArrayBuffer;
bufferSize: number;
};
......@@ -44,8 +45,28 @@ export const readRawContentFromBuffer = (props: {
buffer: Buffer;
}) => {
const bufferSize = props.buffer.length;
const sourceArrayBuffer = props.buffer.buffer;
const canTransferBuffer =
props.buffer.byteOffset === 0 &&
props.buffer.byteLength === sourceArrayBuffer.byteLength &&
sourceArrayBuffer instanceof ArrayBuffer;
if (canTransferBuffer) {
/**
* 大文件解析时优先 transfer 独占 ArrayBuffer,避免再复制一份 SharedArrayBuffer。
* readFile worker 会消费输入 buffer,调用方不应在提交解析后继续复用该 buffer。
*/
return getReadFileWorker().run(
{
extension: props.extension,
encoding: props.encoding,
buffer: sourceArrayBuffer,
bufferSize
},
[sourceArrayBuffer]
);
}
// 使用 SharedArrayBuffer,避免数据复制
const sharedBuffer = new SharedArrayBuffer(bufferSize);
const sharedArray = new Uint8Array(sharedBuffer);
sharedArray.set(props.buffer);
......
import * as pdfjs from 'pdfjs-dist/legacy/build/pdf.mjs';
// @ts-ignore
import('pdfjs-dist/legacy/build/pdf.worker.min.mjs');
import { type ReadRawTextByBuffer, type ReadFileResponse } from '../type';
import { getLogger, LogCategories } from '../../../common/logger';
type TokenType = {
str: string;
dir: string;
width: number;
height: number;
transform: number[];
fontName: string;
hasEOL: boolean;
};
import { type ReadRawTextByBuffer, type ReadFileResponse, type ParsedPage } from '../type';
import { postprocessLiteParsePages } from '../utils/pdfTextPostprocess';
const LITE_PARSE_MAX_PAGES = 100000;
const LITE_PARSE_BATCH_PAGES = 100;
/**
* 使用 LiteParse 解析普通 PDF 文本。
*
* OCR 默认关闭,避免普通系统解析路径产生额外耗时和 OCR 资源依赖;LiteParse 只输出文本
* 和 textItems,本路径不会返回图片或 Markdown。PDF 按页分批解析,降低 PDFium/native
* 一次性持有大量页面结构的内存峰值;最后仍统一后处理全部 pages,保持现有文本清理行为。
*/
export const readPdfFile = async ({ buffer }: ReadRawTextByBuffer): Promise<ReadFileResponse> => {
const logger = getLogger(LogCategories.INFRA.WORKER);
const readPDFPage = async (doc: any, pageNo: number) => {
try {
const page = await doc.getPage(pageNo);
const tokenizedText = await page.getTextContent();
const viewport = page.getViewport({ scale: 1 });
const pageHeight = viewport.height;
const headerThreshold = pageHeight * 0.95;
const footerThreshold = pageHeight * 0.05;
const pageTexts: TokenType[] = tokenizedText.items.filter((token: TokenType) => {
return (
!token.transform ||
(token.transform[5] < headerThreshold && token.transform[5] > footerThreshold)
);
});
const { LiteParse } = await import('@llamaindex/liteparse');
const pages: ParsedPage[] = [];
// concat empty string 'hasEOL'
for (let i = 0; i < pageTexts.length; i++) {
const item = pageTexts[i];
if (item.str === '' && pageTexts[i - 1]) {
pageTexts[i - 1].hasEOL = item.hasEOL;
pageTexts.splice(i, 1);
i--;
}
}
page.cleanup();
return pageTexts
.map((token) => {
const paragraphEnd = token.hasEOL && /([。?!.?!\n\r]|(\r\n))$/.test(token.str);
return paragraphEnd ? `${token.str}\n` : token.str;
})
.join('');
} catch (error) {
logger.error('Failed to read pdf page', { pageNo, error });
return '';
}
};
for (let pageStart = 1; pageStart <= LITE_PARSE_MAX_PAGES; pageStart += LITE_PARSE_BATCH_PAGES) {
const pageEnd = pageStart + LITE_PARSE_BATCH_PAGES - 1;
const parser = new LiteParse({
ocrEnabled: false,
maxPages: LITE_PARSE_MAX_PAGES,
targetPages: `${pageStart}-${pageEnd}`,
quiet: true
});
const result = await parser.parse(buffer);
// Create a completely new ArrayBuffer to avoid SharedArrayBuffer transferList issues
const uint8Array = new Uint8Array(buffer.byteLength);
uint8Array.set(new Uint8Array(buffer.buffer, buffer.byteOffset, buffer.byteLength));
const loadingTask = pdfjs.getDocument({ data: uint8Array });
const doc = await loadingTask.promise;
if (!result.pages.length) break;
const pageArr = Array.from({ length: doc.numPages }, (_, i) => i + 1);
const result = (
await Promise.all(pageArr.map(async (page) => await readPDFPage(doc, page)))
).join('');
pages.push(...result.pages);
}
loadingTask.destroy();
const rawText = postprocessLiteParsePages(pages);
return {
rawText: result
rawText
};
};
......@@ -11,7 +11,8 @@ import { readCsvRawText } from './extension/csv';
type IncomingMessage = {
id: string;
} & Omit<ReadRawTextProps<any>, 'buffer'> & {
sharedBuffer: SharedArrayBuffer;
buffer?: ArrayBuffer;
sharedBuffer?: SharedArrayBuffer;
bufferSize: number;
};
......@@ -40,12 +41,16 @@ const read = async (params: ReadRawTextByBuffer) => {
};
parentPort?.on('message', async (props: IncomingMessage) => {
const { id, sharedBuffer, bufferSize, extension, encoding } = props;
const { id, buffer: transferredBuffer, sharedBuffer, bufferSize, extension, encoding } = props;
try {
// 使用 SharedArrayBuffer,零拷贝共享内存
const sharedArray = new Uint8Array(sharedBuffer);
const buffer = Buffer.from(sharedArray.buffer, 0, bufferSize);
const rawBuffer = transferredBuffer ?? sharedBuffer;
if (!rawBuffer) {
throw new Error('Read file worker missing buffer');
}
// 优先使用 transfer 进来的 ArrayBuffer;兼容旧的 SharedArrayBuffer 零拷贝路径。
const buffer = Buffer.from(rawBuffer, 0, bufferSize);
const data = await read({ extension, encoding, buffer });
......
......@@ -12,6 +12,25 @@ export type ImageType = {
mime: string;
};
export type TextItem = {
text: string;
x: number;
y: number;
width: number;
height: number;
fontName?: string;
fontSize?: number;
confidence?: number;
};
export type ParsedPage = {
pageNum: number;
width: number;
height: number;
text: string;
textItems: TextItem[];
};
export type ReadFileResponse = {
rawText: string;
formatText?: string;
......
import type { ParsedPage, TextItem } from '@llamaindex/liteparse';
const CJK_RE = /[\u3400-\u4dbf\u4e00-\u9fff\uf900-\ufaff]/;
const CJK_END_RE = /[\u3400-\u4dbf\u4e00-\u9fff\uf900-\ufaff)》】」』”]$/;
const CJK_START_RE = /^[\u3400-\u4dbf\u4e00-\u9fff\uf900-\ufaff(《【「『“]/;
const SENTENCE_END_RE = /[。!?!?;;::)】》」』”]$/;
const PARAGRAPH_END_RE = /[。!?!?]$/;
const BULLET_RE = /^(?:[·•●▪-]\s*|\(\d+\)|([一二三四五六七八九十\d]+)|\d+(?:\.\d+)*\s+)/;
const OBVIOUS_HEADING_RE =
/^(?:\s*言|目\s*录|图\s*目\s*录|表\s*目\s*录|参考文献|版权声明|第\s*\d+\s*[章节]|[一二三四五六七八九十]+、|([一二三四五六七八九十]+)|\d+(?:\.\d+)+\s*)/;
const TOC_LINE_RE = /\.{4,}\s*\d+$/;
const PAGE_NO_RE = /^[-—]?\s*\d{1,5}\s*[-—]?$/;
const URL_NOISE_RE = /^\/?[a-z]{2}(?:\/|\))|^\(\/[a-z]{2}\/?\)$/i;
export type PdfTextPostprocessOptions = {
normalizeUnicode?: boolean;
trimPageEdge?: boolean;
headerRatio?: number;
footerRatio?: number;
lineYRatio?: number;
minSpaceGapRatio?: number;
wideSpaceGapRatio?: number;
mergeVisualLines?: boolean;
removeRepeatedPageNoise?: boolean;
repeatedNoiseMinCount?: number;
repeatedNoiseMaxLength?: number;
dropPurePageNumber?: boolean;
inlineNoisePhrases?: string[];
};
type NormalizedTextItem = {
text: string;
x: number;
y: number;
width: number;
height: number;
fontSize: number;
};
type LineGroup = {
y: number;
items: NormalizedTextItem[];
};
const DEFAULT_OPTIONS = {
normalizeUnicode: false,
trimPageEdge: true,
headerRatio: 0.05,
footerRatio: 0.05,
lineYRatio: 0.55,
minSpaceGapRatio: 0.35,
wideSpaceGapRatio: 1.2,
mergeVisualLines: true,
removeRepeatedPageNoise: true,
repeatedNoiseMinCount: 3,
repeatedNoiseMaxLength: 30,
dropPurePageNumber: true,
inlineNoisePhrases: []
} satisfies Required<PdfTextPostprocessOptions>;
/**
* 将 LiteParse 的坐标文本项恢复为更适合知识库切分的纯文本。
*
* 该函数只做低风险文本整理:按 y/x 坐标重组行、过滤页面边缘页眉页脚、
* 删除重复页码/页面噪声,并保守合并 PDF 视觉换行。它不做 OCR、不提取图片,
* 也不把页面截图插入文本,避免改变普通 PDF 解析的成本模型和返回契约。
*/
export const postprocessLiteParsePages = (
pages: Pick<ParsedPage, 'height' | 'textItems'>[],
options: PdfTextPostprocessOptions = {}
) => {
const opts = { ...DEFAULT_OPTIONS, ...options };
const lines = pages.flatMap((page) => extractPageLines(page, opts));
return mergeLines(lines, opts);
};
export const extractPageLines = (
page: Pick<ParsedPage, 'height' | 'textItems'>,
options: PdfTextPostprocessOptions = {}
) => {
const opts = { ...DEFAULT_OPTIONS, ...options };
const items = (page.textItems || [])
.map((item) => normalizeTextItem(item, opts))
.filter((item) => item.text)
.filter((item) => !opts.trimPageEdge || isInsidePageBody(item, page, opts));
if (items.length === 0) return [];
const medianHeight = median(items.map((item) => item.height || item.fontSize || 10)) || 10;
const lineTolerance = Math.max(2, medianHeight * opts.lineYRatio);
const lines: LineGroup[] = [];
for (const item of items.sort((a, b) => a.y - b.y || a.x - b.x)) {
const target = lines.find((line) => Math.abs(line.y - item.y) <= lineTolerance);
if (target) {
target.items.push(item);
target.y = (target.y * (target.items.length - 1) + item.y) / target.items.length;
continue;
}
lines.push({ y: item.y, items: [item] });
}
return lines
.sort((a, b) => a.y - b.y)
.map((line) => joinLineItems(line.items, medianHeight, opts))
.map((line) => line.trim())
.filter(Boolean);
};
const normalizeTextItem = (
item: Partial<TextItem>,
opts: Required<PdfTextPostprocessOptions>
): NormalizedTextItem => {
return {
text: normalizeText(String(item.text || ''), opts).trim(),
x: Number(item.x) || 0,
y: Number(item.y) || 0,
width: Math.max(0, Number(item.width) || 0),
height: Math.max(0, Number(item.height || item.fontSize) || 0),
fontSize: Math.max(0, Number(item.fontSize) || 0)
};
};
const normalizeText = (text: string, opts: Required<PdfTextPostprocessOptions>) => {
const normalized = opts.normalizeUnicode ? text.normalize('NFKC') : text;
return normalized
.replace(/[\u200b\u200c\u200d\ufeff]/g, '')
.replace(/[ \t]+/g, ' ')
.replace(/\s+([,。!?;:、,.!?;:])/g, '$1')
.replace(/([(《【「『“])\s+/g, '$1')
.replace(/\s+([)》】」』”])/g, '$1');
};
const isInsidePageBody = (
item: NormalizedTextItem,
page: Pick<ParsedPage, 'height'>,
opts: Required<PdfTextPostprocessOptions>
) => {
const pageHeight = Number(page.height) || 0;
if (!pageHeight) return true;
const topCutoff = pageHeight * opts.headerRatio;
const bottomCutoff = pageHeight * (1 - opts.footerRatio);
return item.y >= topCutoff && item.y <= bottomCutoff;
};
const joinLineItems = (
items: NormalizedTextItem[],
medianHeight: number,
opts: Required<PdfTextPostprocessOptions>
) => {
let line = '';
let previous: NormalizedTextItem | undefined;
for (const item of items.sort((a, b) => a.x - b.x)) {
if (!previous) {
line = item.text;
previous = item;
continue;
}
const gap = item.x - (previous.x + previous.width);
const shouldSpace =
gap > medianHeight * opts.wideSpaceGapRatio ||
(gap > medianHeight * opts.minSpaceGapRatio && needsSpace(line, item.text));
line = shouldSpace ? `${line} ${item.text}` : joinText(line, item.text);
previous = item;
}
return normalizeText(line, opts);
};
const mergeLines = (lines: string[], opts: Required<PdfTextPostprocessOptions>) => {
const noiseSet = opts.removeRepeatedPageNoise
? findRepeatedNoise(lines, opts)
: new Set<string>();
const paragraphs: string[] = [];
let current = '';
let previousStandalone = false;
const flush = () => {
if (current) paragraphs.push(current);
current = '';
};
for (const rawLine of lines) {
const line = cleanupInlineNoise(normalizeText(rawLine, opts).trim(), noiseSet, opts);
if (!line) continue;
if (noiseSet.has(line)) continue;
if (opts.dropPurePageNumber && PAGE_NO_RE.test(line)) continue;
const standalone = isStandaloneLine(line);
if (!current) {
current = line;
previousStandalone = standalone;
if (standalone) flush();
continue;
}
if (standalone) {
flush();
paragraphs.push(line);
previousStandalone = true;
continue;
}
if (shouldMergeLine(current, line, previousStandalone, opts)) {
current = joinText(current, line);
} else {
flush();
current = line;
}
previousStandalone = false;
}
flush();
return paragraphs.join('\n\n') + (paragraphs.length > 0 ? '\n' : '');
};
const shouldMergeLine = (
current: string,
next: string,
previousStandalone: boolean,
opts: Required<PdfTextPostprocessOptions>
) => {
if (!opts.mergeVisualLines) return false;
if (previousStandalone) return false;
if (SENTENCE_END_RE.test(current)) return false;
if (isStandaloneLine(next)) return false;
return true;
};
const isStandaloneLine = (line: string) => {
if (PAGE_NO_RE.test(line)) return true;
if (OBVIOUS_HEADING_RE.test(line)) return true;
if (TOC_LINE_RE.test(line)) return true;
if (BULLET_RE.test(line)) return true;
if (URL_NOISE_RE.test(line)) return true;
if (line.length <= 14 && CJK_RE.test(line) && !/[,,。!?;;::]/.test(line)) return true;
return false;
};
const joinText = (left: string, right: string) => {
if (!left) return right;
if (!right) return left;
if (needsSpace(left, right)) return `${left} ${right}`;
return `${left}${right}`;
};
const needsSpace = (left: string, right: string) => {
if (!left || !right) return false;
if (CJK_END_RE.test(left) && CJK_START_RE.test(right)) return false;
if (/[-/([{]$/.test(left)) return false;
if (/^[,.;:!?%)}\]]/.test(right)) return false;
return /[A-Za-z0-9]$/.test(left) || /^[A-Za-z0-9]/.test(right);
};
const findRepeatedNoise = (lines: string[], opts: Required<PdfTextPostprocessOptions>) => {
const counts = new Map<string, number>();
for (const line of lines) {
const text = normalizeText(line, opts).trim();
if (!isNoiseCandidate(text, opts)) continue;
counts.set(text, (counts.get(text) || 0) + 1);
}
return new Set(
[...counts.entries()]
.filter(([, count]) => count >= opts.repeatedNoiseMinCount)
.map(([text]) => text)
);
};
const isNoiseCandidate = (line: string, opts: Required<PdfTextPostprocessOptions>) => {
if (!line) return false;
if (PAGE_NO_RE.test(line)) return true;
if (URL_NOISE_RE.test(line)) return true;
if (line.length > opts.repeatedNoiseMaxLength) return false;
if (!SENTENCE_END_RE.test(line) && !PARAGRAPH_END_RE.test(line)) return true;
return false;
};
const cleanupInlineNoise = (
line: string,
noiseSet: Set<string>,
opts: Required<PdfTextPostprocessOptions>
) => {
let text = line;
const candidates = new Set(opts.inlineNoisePhrases);
// 只对 URL/语言路径类噪声做行内删除,避免把短词误删出正文。
for (const noise of noiseSet) {
if (URL_NOISE_RE.test(noise)) candidates.add(noise);
}
for (const noise of candidates) {
if (!noise) continue;
text = text.split(noise).join(' ');
}
return normalizeText(text, opts).trim();
};
const median = (nums: number[]) => {
const valid = nums.filter(Number.isFinite).sort((a, b) => a - b);
if (valid.length === 0) return 0;
return valid[Math.floor(valid.length / 2)];
};
import type { Worker as NodeWorker } from 'worker_threads';
import type { TransferListItem, Worker as NodeWorker } from 'worker_threads';
import path from 'path';
import { getLogger, LogCategories } from '../common/logger';
import { serviceEnv } from '../env';
......@@ -65,7 +65,12 @@ export const runWorker = <T = any>(name: WorkerNameEnum, params?: Record<string,
});
};
type WorkerRunTaskType<T> = { data: T; resolve: (e: any) => void; reject: (e: any) => void };
type WorkerRunTaskType<T> = {
data: T;
transferList?: TransferListItem[];
resolve: (e: any) => void;
reject: (e: any) => void;
};
type WorkerQueueItem = {
id: string;
worker: NodeWorker;
......@@ -116,7 +121,7 @@ export class WorkerPool<Props = Record<string, any>, Response = any> {
this.maxTasksPerWorker = maxTasksPerWorker;
}
private runTask({ data, resolve, reject }: WorkerRunTaskType<Props>) {
private runTask({ data, transferList, resolve, reject }: WorkerRunTaskType<Props>) {
// Get idle worker or create a new worker
const runningWorker = (() => {
// @ts-ignore
......@@ -145,23 +150,27 @@ export class WorkerPool<Props = Record<string, any>, Response = any> {
this.deleteWorker(runningWorker.id);
}, this.taskTimeoutMs);
runningWorker.worker.postMessage({
id: runningWorker.id,
...data
});
runningWorker.worker.postMessage(
{
id: runningWorker.id,
...data
},
transferList
);
} else {
// Not enough worker, push to wait queue
this.waitQueue.push({ data, resolve, reject });
this.waitQueue.push({ data, transferList, resolve, reject });
}
}
run(data: Props) {
run(data: Props, transferList?: TransferListItem[]) {
return new Promise<Response>((resolve, reject) => {
/*
Whether the task is executed immediately or delayed, the promise callback will dispatch after task complete.
*/
this.runTask({
data,
transferList,
resolve,
reject
});
......
......@@ -33,6 +33,9 @@ catalogs:
'@emotion/styled':
specifier: ^11
version: 11.11.0
'@llamaindex/liteparse':
specifier: 2.0.5
version: 2.0.5
'@modelcontextprotocol/sdk':
specifier: ^1
version: 1.26.0
......@@ -464,6 +467,9 @@ importers:
'@fastgpt/global':
specifier: workspace:*
version: link:../global
'@llamaindex/liteparse':
specifier: 'catalog:'
version: 2.0.5
'@mariozechner/pi-agent-core':
specifier: ^0.67.3
version: 0.67.3(@modelcontextprotocol/sdk@1.26.0(zod@4.1.12))(bufferutil@4.1.0)(utf-8-validate@5.0.10)(ws@8.20.0(bufferutil@4.1.0)(utf-8-validate@5.0.10))(zod@4.1.12)
......@@ -605,9 +611,6 @@ importers:
papaparse:
specifier: 5.4.1
version: 5.4.1
pdfjs-dist:
specifier: 4.10.38
version: 4.10.38
pg:
specifier: ^8.10.0
version: 8.14.0
......@@ -1274,6 +1277,9 @@ importers:
'@fortaine/fetch-event-source':
specifier: ^3.0.6
version: 3.0.6
'@llamaindex/liteparse':
specifier: 'catalog:'
version: 2.0.5
'@modelcontextprotocol/sdk':
specifier: 'catalog:'
version: 1.26.0(zod@4.1.12)
......@@ -4100,6 +4106,49 @@ packages:
'@lezer/yaml@1.0.3':
resolution: {integrity: sha512-GuBLekbw9jDBDhGur82nuwkxKQ+a3W5H0GfaAthDXcAu+XdpS43VlnxA9E9hllkpSP5ellRDKjLLj7Lu9Wr6xA==}
'@llamaindex/liteparse-darwin-arm64@2.0.5':
resolution: {integrity: sha512-Z6gO+2Anwx7baBDXNUMHl8QYWuaDmLQZCZ4NKBgMXGlBaESL6OWsShigFMtPWRlNw5EsP/r+tn7uJmPxy2z2tg==}
cpu: [arm64]
os: [darwin]
'@llamaindex/liteparse-darwin-x64@2.0.5':
resolution: {integrity: sha512-vK+BWVPrMId90sTCg5D/5I7ei2DLuy9Zt0mmByphKAGhiEb906X/hQJELKf4a6pBm1QhS0rb8sPDoHhQTY9aEA==}
cpu: [x64]
os: [darwin]
'@llamaindex/liteparse-linux-arm64-gnu@2.0.5':
resolution: {integrity: sha512-CKClEmHFf1iEwMp7SHqfcpa12D1Esr+JkxxE3pWS0fTXoTEkDOrjRQP1AEbXczGqb+pVVhCxWYcp5DsiXwC5aQ==}
cpu: [arm64]
os: [linux]
libc: [glibc]
'@llamaindex/liteparse-linux-x64-gnu@2.0.5':
resolution: {integrity: sha512-3FBSjROuqleSLOhXFPEoG6Rq0mFGVxGuN45d1x3fxqUfDig753rTdAqisJnyLYLDiEUjK9llPP3rmuuFBVBFBg==}
cpu: [x64]
os: [linux]
libc: [glibc]
'@llamaindex/liteparse-linux-x64-musl@2.0.5':
resolution: {integrity: sha512-arucwQ/pqetqtno+wngn3ho+oof7zJwCaNkJ63YaTk9xyxxWUQ9h8Dl4IdI1tCprgPx0U9DOWtKfJzcASXteUA==}
cpu: [x64]
os: [linux]
libc: [musl]
'@llamaindex/liteparse-win32-arm64-msvc@2.0.5':
resolution: {integrity: sha512-0yT8qm7tkvZv/LW/zS8WuWA7ixLmkuFI9DjSO1dcIZpq/RgyvP/jyZsNSZEqR3IfR1GsViFmplTp9MjW7KV+1Q==}
cpu: [arm64]
os: [win32]
'@llamaindex/liteparse-win32-x64-msvc@2.0.5':
resolution: {integrity: sha512-htxyLATNuWJdTVHCM73jkDBNJZrvzz7gptUpGsVjhYOc7WvTBNCPsoSeOo07NonBH4mwEZfOU3XjuEwwjYGDWQ==}
cpu: [x64]
os: [win32]
'@llamaindex/liteparse@2.0.5':
resolution: {integrity: sha512-EQN2hYKRSDqhl5SAr8kqJZFn3pzh6DXAhhtng2zRlqvg/rslaroaBAbvzKep+ElljrmfYNZYjJqavgDK1IM1TQ==}
engines: {node: '>=18.0.0'}
hasBin: true
'@logtape/logtape@2.0.2':
resolution: {integrity: sha512-cveUBLbCMFkvkLycP/2vNWvywl47JcJbazHIju94/QNGboZ/jyYAgZIm0ZXezAOx3eIz8OG1EOZ5CuQP3+2FQg==}
......@@ -7903,6 +7952,10 @@ packages:
resolution: {integrity: sha512-yPVavfyCcRhmorC7rWlkHn15b4wDVgVmBA7kV4QVBsF7kv/9TKJAbAXVTxvTnwP8HHKjRCJDClKbciiYS7p0DQ==}
engines: {node: '>=16'}
commander@12.1.0:
resolution: {integrity: sha512-Vw8qHK3bZM9y/P10u3Vib8o/DdkvA2OtPtZvD871QKjy74Wj1WSKFILMPRPSdUSx5RFK1arlJzEtA4PkFgnbuA==}
engines: {node: '>=18'}
commander@14.0.3:
resolution: {integrity: sha512-H+y0Jo/T1RZ9qPP4Eh1pkcQcLRglraJaSLoyOtHxu6AapkjWVCy2Sit1QQ4x3Dng8qDlSsZEet7g5Pq06MvTgw==}
engines: {node: '>=20'}
......@@ -12018,10 +12071,6 @@ packages:
pause-stream@0.0.11:
resolution: {integrity: sha512-e3FBlXLmN/D1S+zHzanP4E/4Z60oFAa3O051qt1pxa7DEJWKAyil6upYVXCWadEnuoqa4Pkc9oUx9zsxYeRv8A==}
pdfjs-dist@4.10.38:
resolution: {integrity: sha512-/Y3fcFrXEAsMjJXeL9J8+ZG9U01LbuWaYypvDW2ycW1jL269L3js3DVBjDJ0Up9Np1uqDXsDrRihHANhZOlwdQ==}
engines: {node: '>=20'}
pend@1.2.0:
resolution: {integrity: sha512-F3asv42UuXchdzt+xXqfW1OGlVBe+mxa2mqI0pg5yAHZPvFmY3Y6drSf/GQ1A86WgWEN9Kzh/WrgKa6iGcHXLg==}
......@@ -18061,6 +18110,39 @@ snapshots:
'@lezer/highlight': 1.2.1
'@lezer/lr': 1.4.2
'@llamaindex/liteparse-darwin-arm64@2.0.5':
optional: true
'@llamaindex/liteparse-darwin-x64@2.0.5':
optional: true
'@llamaindex/liteparse-linux-arm64-gnu@2.0.5':
optional: true
'@llamaindex/liteparse-linux-x64-gnu@2.0.5':
optional: true
'@llamaindex/liteparse-linux-x64-musl@2.0.5':
optional: true
'@llamaindex/liteparse-win32-arm64-msvc@2.0.5':
optional: true
'@llamaindex/liteparse-win32-x64-msvc@2.0.5':
optional: true
'@llamaindex/liteparse@2.0.5':
dependencies:
commander: 12.1.0
optionalDependencies:
'@llamaindex/liteparse-darwin-arm64': 2.0.5
'@llamaindex/liteparse-darwin-x64': 2.0.5
'@llamaindex/liteparse-linux-arm64-gnu': 2.0.5
'@llamaindex/liteparse-linux-x64-gnu': 2.0.5
'@llamaindex/liteparse-linux-x64-musl': 2.0.5
'@llamaindex/liteparse-win32-arm64-msvc': 2.0.5
'@llamaindex/liteparse-win32-x64-msvc': 2.0.5
'@logtape/logtape@2.0.2': {}
'@logtape/pretty@2.0.2(@logtape/logtape@2.0.2)':
......@@ -22590,6 +22672,8 @@ snapshots:
commander@11.1.0: {}
commander@12.1.0: {}
commander@14.0.3: {}
commander@2.20.3: {}
......@@ -27864,10 +27948,6 @@ snapshots:
dependencies:
through: 2.3.8
pdfjs-dist@4.10.38:
optionalDependencies:
'@napi-rs/canvas': 0.1.69
pend@1.2.0: {}
performance-now@2.1.0: {}
......@@ -25,6 +25,7 @@ catalog:
'@dnd-kit/core': ^6.3.1
'@emotion/react': ^11
'@emotion/styled': ^11
'@llamaindex/liteparse': 2.0.5
'@modelcontextprotocol/sdk': ^1
'@node-rs/jieba': 2.0.1
'@svgr/webpack': ^6.5.1
......
......@@ -69,7 +69,7 @@ RUN addgroup --system --gid 1001 nodejs
RUN adduser --system --uid 1001 nextjs
RUN [ -z "$proxy" ] || sed -i 's/dl-cdn.alpinelinux.org/mirrors.ustc.edu.cn/g' /etc/apk/repositories
RUN apk add --no-cache curl ca-certificates \
RUN apk add --no-cache curl ca-certificates libc++ \
&& update-ca-certificates
# copy running files
......
......@@ -44,6 +44,7 @@
"@fastgpt/service": "workspace:*",
"@fastgpt/web": "workspace:*",
"@fortaine/fetch-event-source": "^3.0.6",
"@llamaindex/liteparse": "catalog:",
"@modelcontextprotocol/sdk": "catalog:",
"@monaco-editor/react": "^4.7.0",
"@node-rs/jieba": "catalog:",
......
import { build, BuildOptions, context } from 'esbuild';
import fs from 'fs';
import { createRequire } from 'module';
import path from 'path';
// 项目路径
const ROOT_DIR = path.resolve(__dirname, '../../..');
const WORKER_SOURCE_DIR = path.join(ROOT_DIR, 'packages/service/worker');
const WORKER_OUTPUT_DIR = path.join(__dirname, '../worker');
const WORKER_RUNTIME_NODE_MODULES_DIR = path.join(WORKER_OUTPUT_DIR, 'node_modules');
const OTEL_SDK_DIR = path.join(ROOT_DIR, 'sdk/otel/src');
const require = createRequire(import.meta.url);
const workerRuntimePackages = ['@llamaindex/liteparse'];
const liteParsePlatformPackages = [
'@llamaindex/liteparse-darwin-x64',
'@llamaindex/liteparse-darwin-arm64',
'@llamaindex/liteparse-linux-x64-gnu',
'@llamaindex/liteparse-linux-x64-musl',
'@llamaindex/liteparse-linux-arm64-gnu',
'@llamaindex/liteparse-linux-arm64-musl',
'@llamaindex/liteparse-win32-x64-msvc',
'@llamaindex/liteparse-win32-arm64-msvc'
];
const resolvePackageDir = (packageName: string, resolvePaths: string[]) => {
try {
return path.dirname(require.resolve(`${packageName}/package.json`, { paths: resolvePaths }));
} catch {
return;
}
};
const copyPackage = (packageName: string, sourceDir: string) => {
const destination = path.join(WORKER_RUNTIME_NODE_MODULES_DIR, ...packageName.split('/'));
fs.rmSync(destination, { recursive: true, force: true });
fs.mkdirSync(path.dirname(destination), { recursive: true });
fs.cpSync(sourceDir, destination, {
recursive: true,
dereference: true
});
console.log(`📦 ${packageName} 运行时依赖已复制 → ${path.relative(process.cwd(), destination)}`);
};
/**
* 复制 worker external 依赖到 worker 目录下。
*
* 这些依赖不适合直接打进 esbuild bundle:LiteParse 包含 N-API .node 文件和 PDFium
* 动态库,必须以真实文件形式存在。Docker runner 已经复制整个 worker 目录,因此把
* runtime node_modules 放在 worker 旁边即可让 Node worker 线程就近解析。
*/
const copyWorkerRuntimePackages = () => {
fs.rmSync(WORKER_RUNTIME_NODE_MODULES_DIR, { recursive: true, force: true });
for (const packageName of workerRuntimePackages) {
const sourceDir = resolvePackageDir(packageName, [__dirname, ROOT_DIR]);
if (!sourceDir) {
throw new Error(`Worker runtime dependency "${packageName}" is not installed.`);
}
copyPackage(packageName, sourceDir);
if (packageName === '@llamaindex/liteparse') {
for (const platformPackage of liteParsePlatformPackages) {
const platformSourceDir = resolvePackageDir(platformPackage, [sourceDir, __dirname, ROOT_DIR]);
if (!platformSourceDir) continue;
copyPackage(platformPackage, platformSourceDir);
}
}
}
};
/**
* Worker 预编译脚本
......@@ -54,6 +119,7 @@ async function buildWorkers(watch: boolean = false) {
'@fastgpt-sdk/otel/metrics': path.join(OTEL_SDK_DIR, 'metrics-entry.ts'),
'@fastgpt-sdk/otel/tracing': path.join(OTEL_SDK_DIR, 'tracing-entry.ts')
},
external: ['@llamaindex/liteparse', '@llamaindex/liteparse-*'],
// 移除调试代码
drop: process.env.NODE_ENV === 'production' ? ['console', 'debugger'] : []
};
......@@ -87,6 +153,8 @@ async function buildWorkers(watch: boolean = false) {
// 过滤掉失败的 context
const validContexts = contexts.filter((ctx) => ctx !== null);
copyWorkerRuntimePackages();
console.log('━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━');
console.log(`✅ ${validContexts.length}/${workers.length} 个 Worker 正在监听中`);
console.log('━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━');
......@@ -137,6 +205,7 @@ async function buildWorkers(watch: boolean = false) {
// 非监听模式下,如果有失败的编译,退出并返回错误码
process.exit(1);
}
copyWorkerRuntimePackages();
console.log('━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━');
}
}
......
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