From 569bb5c47ceb8acbb096f611eb0b466caf5b90bd Mon Sep 17 00:00:00 2001 From: FinleyGe Date: Wed, 25 Jun 2025 20:48:07 +0800 Subject: [PATCH] chore: type --- src/worker/index.ts | 45 ++++++++++-------- src/worker/type.ts | 48 +++++++++++++++++++ src/worker/utils.ts | 20 ++++++++ src/worker/worker.d.ts | 21 +++++++++ src/worker/worker.ts | 104 +++++++++++++++++++++++++++-------------- tsconfig.json | 3 +- 6 files changed, 187 insertions(+), 54 deletions(-) create mode 100644 src/worker/type.ts create mode 100644 src/worker/utils.ts create mode 100644 src/worker/worker.d.ts diff --git a/src/worker/index.ts b/src/worker/index.ts index e5f99984..4def47d3 100644 --- a/src/worker/index.ts +++ b/src/worker/index.ts @@ -4,6 +4,7 @@ import { ToolCallbackReturnSchema } from '../../packages/tool/type/tool'; import { z } from 'zod'; import { addLog } from '@/utils/log'; import { isProd } from '@/constants'; +import type { Worker2MainMessageType } from './type'; type WorkerQueueItem = { id: string; @@ -183,24 +184,27 @@ export async function dispatchWithNewWorker(data: { const resolvePromise = new Promise>( (resolve, reject) => { - worker.on( - 'message', - ({ type, data }: WorkerResponse>) => { - if (type === 'success') { - resolve(data); - worker.terminate(); - } else if (type === 'error') { - reject(data); - worker.terminate(); - } else if (type === 'log') { - const msg = data as { - type: 'info' | 'error' | 'warn'; - args: any[]; - }; - addLog[msg.type](`Tool run: `, msg.args); - } + worker.on('message', ({ type, data }: Worker2MainMessageType) => { + if (type === 'success') { + resolve(data); + worker.terminate(); + } else if (type === 'error') { + reject(data); + worker.terminate(); + } else if (type === 'log') { + const msg = data as { + type: 'info' | 'error' | 'warn'; + args: any[]; + }; + addLog[msg.type](`Tool run: `, msg.args); + } else if (type === 'uploadFile') { + // TODO upload + worker.postMessage({ + type: 'uploadFileResponse', + data: null // TODO: response + }); } - ); + }); worker.on('error', (err) => { addLog.error(`Run tool error`, err); @@ -214,8 +218,11 @@ export async function dispatchWithNewWorker(data: { }); worker.postMessage({ - toolDirName: tool.toolDirName, - ...data + type: 'runTool', + data: { + toolDirName: tool.toolDirName, + ...data + } }); } ); diff --git a/src/worker/type.ts b/src/worker/type.ts new file mode 100644 index 00000000..de722cd0 --- /dev/null +++ b/src/worker/type.ts @@ -0,0 +1,48 @@ +import z from 'zod'; + +/** + * Worker --> Main Thread + */ +export const Worker2MainMessageSchema = z.discriminatedUnion('type', [ + z.object({ + type: z.literal('uploadFile'), + data: z.any() // minio upload file params + }), + z.object({ + type: z.literal('log'), + data: z.object({ + type: z.enum(['info', 'error', 'warn']), + args: z.array(z.any()) + }) + }), + z.object({ + type: z.literal('success'), + data: z.any() + }), + z.object({ + type: z.literal('error'), + data: z.any() + }) +]); + +/** + * Main Thread --> Worker + */ +export const Main2WorkerMessageSchema = z.discriminatedUnion('type', [ + z.object({ + type: z.literal('runTool'), + data: z.object({ + toolId: z.string(), + inputs: z.any(), + systemVar: z.any(), + toolDirName: z.string() + }) + }), + z.object({ + type: z.literal('uploadFileResponse'), + data: z.any() // minio upload file response + }) +]); + +export type Worker2MainMessageType = z.infer; +export type Main2WorkerMessageType = z.infer; diff --git a/src/worker/utils.ts b/src/worker/utils.ts new file mode 100644 index 00000000..825b6207 --- /dev/null +++ b/src/worker/utils.ts @@ -0,0 +1,20 @@ +import { parentPort } from 'worker_threads'; + +export const uploadFile = async (data: any) => { + return new Promise((resolve, reject) => { + global.uploadFileResponseFn = (res: any) => { + resolve(res); + }; + parentPort?.postMessage({ + type: 'uploadFile', + data + }); + }); +}; + +declare global { + // eslint-disable-next-line no-var + var uploadFileResponseFn: (data: any) => void | undefined; +} + +export {}; diff --git a/src/worker/worker.d.ts b/src/worker/worker.d.ts new file mode 100644 index 00000000..1d254e6d --- /dev/null +++ b/src/worker/worker.d.ts @@ -0,0 +1,21 @@ +// declare global { +// // eslint-disable-next-line no-var +// var uploadFileResponseFn: ( +// data: Record +// ) => ((data: Record) => void) | undefined; +// } + +//eslint-disable-next-line no-var +declare var uploadFileResponseFn: ( + data: Record +) => ((data: Record) => void) | undefined; + +// declare module NodeJS { +// interface Global { +// uploadFileResponseFn: ( +// data: Record +// ) => ((data: Record) => void) | undefined; +// } +// } + +// export {}; diff --git a/src/worker/worker.ts b/src/worker/worker.ts index da05a8c7..8ead926e 100644 --- a/src/worker/worker.ts +++ b/src/worker/worker.ts @@ -2,8 +2,8 @@ import { parentPort } from 'worker_threads'; import path from 'path'; import { LoadToolsByFilename } from '@tool/init'; import { isProd } from '@/constants'; -import type { SystemVarType } from '@tool/type'; import { getErrText } from '@tool/utils/err'; +import type { Main2WorkerMessageType } from './type'; // rewrite console.debug to send to parent console.debug = (...args: any[]) => { @@ -52,40 +52,76 @@ const basePath = isProd ? process.env.TOOLS_DIR || path.join(process.cwd(), 'dist', 'tools') : path.join(process.cwd(), 'packages', 'tool', 'packages'); -parentPort?.on( - 'message', - async ({ - toolId, - inputs, - systemVar, - toolDirName - }: { - toolId: string; - inputs: Record; - systemVar: SystemVarType; - toolDirName: string; - }) => { - const tools = await LoadToolsByFilename(basePath, toolDirName); - const tool = tools.find((tool) => tool.toolId === toolId); +parentPort?.on('message', async (params: Main2WorkerMessageType) => { + const { type, data } = params; + switch (type) { + case 'runTool': { + const tools = await LoadToolsByFilename(basePath, data.toolDirName); + const tool = tools.find((tool) => tool.toolId === data.toolId); - if (!tool || !tool.cb) { - parentPort?.postMessage({ - type: 'error', - data: `Tool with ID ${toolId} not found or does not have a callback.` - }); - } - try { - const result = await tool?.cb(inputs, systemVar); + if (!tool || !tool.cb) { + parentPort?.postMessage({ + type: 'error', + data: `Tool with ID ${data.toolId} not found or does not have a callback.` + }); + } + try { + const result = tool?.cb(data.inputs, data.systemVar); - parentPort?.postMessage({ - type: 'success', - data: result - }); - } catch (error) { - parentPort?.postMessage({ - type: 'error', - data: getErrText(error) - }); + parentPort?.postMessage({ + type: 'success', + data: result + }); + } catch (error) { + // TODO: 处理错误 + parentPort?.postMessage({ + type: 'error', + data: getErrText(error) + }); + } + break; + } + case 'uploadFileResponse': { + global.uploadFileResponseFn?.(data); + break; } } -); +}); + +// parentPort?.on( +// 'message', +// async ({ +// toolId, +// inputs, +// systemVar, +// toolDirName +// }: { +// toolId: string; +// inputs: Record; +// systemVar: SystemVarType; +// toolDirName: string; +// }) => { +// const tools = await LoadToolsByFilename(basePath, toolDirName); +// const tool = tools.find((tool) => tool.toolId === toolId); + +// if (!tool || !tool.cb) { +// parentPort?.postMessage({ +// type: 'error', +// data: `Tool with ID ${toolId} not found or does not have a callback.` +// }); +// } +// try { +// const result = await tool?.cb(inputs, systemVar); + +// parentPort?.postMessage({ +// type: 'success', +// data: result +// }); +// } catch (error) { +// parentPort?.postMessage({ +// type: 'error', +// data: getErrText(error) +// }); +// } +// } +// ); diff --git a/tsconfig.json b/tsconfig.json index ee28d14f..df8baa63 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -30,5 +30,6 @@ "@/*": ["./src/*"], "@tool/*": ["./packages/tool/*"] } - } + }, + "include": ["**/*.ts", "**/*.d.ts"] }