From 1c6659def5795ae68cab90865d667c6281024544 Mon Sep 17 00:00:00 2001 From: sunlei Date: Thu, 18 Jun 2026 03:54:29 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=A2=9E=E5=8A=A0=E9=80=9A=E7=94=A8?= =?UTF-8?q?=E6=8F=92=E4=BB=B6worker=E8=BF=90=E8=A1=8C=E6=97=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../package/plugin-package-source.service.ts | 22 +- .../integration/runtime/index.ts | 3 + .../runtime/plugin-worker-runtime.factory.ts | 304 +++++++++ .../runtime/plugin-worker.thread.ts | 640 ++++++++++++++++++ .../runtime/worker-runtime.types.ts | 1 + .../plugin-lifecycle-runtime.spec.ts | 32 + .../plugin-worker-entry.spec.ts | 114 ++++ 7 files changed, 1113 insertions(+), 3 deletions(-) create mode 100644 src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker-runtime.factory.ts create mode 100644 src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker.thread.ts create mode 100644 test/modules/qqbot/plugin-platform/plugin-worker-entry.spec.ts diff --git a/src/modules/qqbot/plugin-platform/infrastructure/integration/package/plugin-package-source.service.ts b/src/modules/qqbot/plugin-platform/infrastructure/integration/package/plugin-package-source.service.ts index 8aa8187..ea8997f 100644 --- a/src/modules/qqbot/plugin-platform/infrastructure/integration/package/plugin-package-source.service.ts +++ b/src/modules/qqbot/plugin-platform/infrastructure/integration/package/plugin-package-source.service.ts @@ -16,7 +16,9 @@ export class QqbotPluginPackageSourceService { * Creates a manifest-only package descriptor source. * @param pathPolicy - Root and entry policy that keeps package discovery inside controlled directories. */ - constructor(private readonly pathPolicy: QqbotPluginPackagePathPolicyService) {} + constructor( + private readonly pathPolicy: QqbotPluginPackagePathPolicyService, + ) {} /** * Discovers package descriptors from one-level package directories under controlled roots. @@ -54,6 +56,21 @@ export class QqbotPluginPackageSourceService { } const manifestLike = JSON.parse(readFileSync(manifestFile, 'utf8')); + return this.resolveDescriptor(controlledPackageRoot, manifestLike); + } + + /** + * Resolves a descriptor from an already available manifest without re-reading `plugin.json`. + * @param packageRoot - Installed or discovered QQBot plugin package directory. + * @param manifestLike - Manifest JSON stored for the package version or read from disk. + * @returns Descriptor with a policy-normalized package root and entry file. + */ + resolveDescriptor( + packageRoot: string, + manifestLike: unknown, + ): QqbotPluginPackageDescriptor { + const controlledPackageRoot = + this.pathPolicy.assertControlledPackageRoot(packageRoot); const manifest = parseQqbotPluginManifest(manifestLike, { pluginRoot: controlledPackageRoot, }); @@ -61,14 +78,13 @@ export class QqbotPluginPackageSourceService { controlledPackageRoot, manifest.entry, ); - const pluginKey = manifest.key; return { entry: manifest.entry, entryFile, manifest, packageRoot: controlledPackageRoot, - pluginKey, + pluginKey: manifest.key, }; } diff --git a/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/index.ts b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/index.ts index 5ea93c2..6f2901d 100644 --- a/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/index.ts +++ b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/index.ts @@ -2,3 +2,6 @@ export * from './worker-runtime'; export * from './worker-runtime.types'; export * from './builtin-plugin-worker-runtime.factory'; export * from './bullmq-plugin-worker-request.queue'; +export * from './plugin-host-bridge.service'; +export * from './plugin-host-bridge.types'; +export * from './plugin-worker-runtime.factory'; diff --git a/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker-runtime.factory.ts b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker-runtime.factory.ts new file mode 100644 index 0000000..f6d4add --- /dev/null +++ b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker-runtime.factory.ts @@ -0,0 +1,304 @@ +import { Injectable } from '@nestjs/common'; +import { ConfigService } from '@nestjs/config'; +import { join } from 'node:path'; +import { Worker } from 'node:worker_threads'; + +import { QqbotPluginPackageSourceService } from '../package/plugin-package-source.service'; +import type { + QqbotPluginPackageDescriptor, + QqbotPluginRuntimeConfigSnapshot, +} from '../package/plugin-package.types'; +import { QqbotPluginHostBridgeService } from './plugin-host-bridge.service'; +import { + createQqbotBullmqWorkerQueueOptions, + QqbotBullmqPluginWorkerRequestQueue, +} from './bullmq-plugin-worker-request.queue'; +import { + QqbotPluginWorkerResponseError, + QqbotPluginWorkerRuntime, + serializePluginWorkerResponseError, +} from './worker-runtime'; +import type { + QqbotPluginWorkerDriver, + QqbotPluginWorkerRequest, +} from './worker-runtime.types'; +import type { QqbotPluginRuntimeFactory } from '@/modules/qqbot/plugin-platform/application/plugin-platform.service'; +import type { + QqbotPluginInstallation, + QqbotPluginVersion, +} from '@/modules/qqbot/plugin-platform/infrastructure/persistence'; + +type WorkerBridgeMessage = + | { + error?: { message?: string; name?: string; stack?: string }; + ok: boolean; + requestId: string; + result?: unknown; + type: 'response'; + } + | { + args?: Record; + method: string; + requestId: string; + type: 'hostCall'; + }; + +type PendingRequest = { + reject: (reason?: unknown) => void; + resolve: (value: unknown) => void; +}; + +type QqbotPluginWorkerThreadDriverOptions = { + configSnapshot: QqbotPluginRuntimeConfigSnapshot; + descriptor: QqbotPluginPackageDescriptor; + installationId: string; + pluginKey: string; +}; + +@Injectable() +export class QqbotPluginWorkerRuntimeFactoryService implements QqbotPluginRuntimeFactory { + /** + * Creates the descriptor-based worker runtime factory. + * @param configService - Nest config source used for BullMQ queue options and runtime config snapshots. + * @param packageSource - Package descriptor resolver that applies controlled-root and entry policies. + * @param hostBridge - Generic host bridge used by worker host-call RPC messages. + */ + constructor( + private readonly configService: ConfigService, + private readonly packageSource: QqbotPluginPackageSourceService, + private readonly hostBridge: QqbotPluginHostBridgeService, + ) {} + + /** + * Creates a worker runtime for one plugin installation and version descriptor. + * @param installation - Installation row whose `installedPath` identifies the package root. + * @param version - Version row whose `manifestJson` is the source of runtime manifest semantics. + * @returns Worker runtime backed by a descriptor-based worker thread and BullMQ serialization queue. + */ + create( + installation: QqbotPluginInstallation, + version: QqbotPluginVersion, + ): QqbotPluginWorkerRuntime { + const descriptor = this.packageSource.resolveDescriptor( + installation.installedPath, + version.manifestJson, + ); + const configSnapshot = this.createConfigSnapshot(descriptor); + const driver = new QqbotPluginWorkerThreadDriver(this.hostBridge, { + configSnapshot, + descriptor, + installationId: installation.id, + pluginKey: descriptor.pluginKey, + }); + + return new QqbotPluginWorkerRuntime( + new QqbotBullmqPluginWorkerRequestQueue( + driver, + createQqbotBullmqWorkerQueueOptions( + this.configService, + descriptor.pluginKey, + installation.id, + ), + ), + { + configSnapshot, + defaultTimeoutMs: descriptor.manifest.runtime.timeoutMs, + descriptor, + installationId: installation.id, + pluginKey: descriptor.pluginKey, + }, + ); + } + + /** + * Captures manifest-declared runtime config keys without hard-coding plugin-specific names. + * @param descriptor - Package descriptor whose manifest owns the config key declarations. + * @returns String snapshot preserving missing keys as `undefined`. + */ + private createConfigSnapshot( + descriptor: QqbotPluginPackageDescriptor, + ): QqbotPluginRuntimeConfigSnapshot { + return Object.fromEntries( + descriptor.manifest.runtime.configKeys.map((key) => { + const value = this.configService.get( + key, + ); + return [ + key, + value === undefined || value === null ? undefined : `${value}`, + ]; + }), + ); + } +} + +export class QqbotPluginWorkerThreadDriver implements QqbotPluginWorkerDriver { + private readonly pendingRequests = new Map(); + private worker?: Worker; + + /** + * Creates a thread driver for one descriptor-based plugin runtime. + * @param hostBridge - Generic platform host bridge used to satisfy worker host calls. + * @param options - Descriptor, installation id, plugin key, and config snapshot passed to workerData. + */ + constructor( + private readonly hostBridge: QqbotPluginHostBridgeService, + private readonly options: QqbotPluginWorkerThreadDriverOptions, + ) {} + + /** + * Sends one runtime request to the worker thread and waits for its response. + * @param message - Lifecycle, operation, task, or event request produced by QqbotPluginWorkerRuntime. + * @returns Worker response payload for the requested plugin action. + */ + async request(message: QqbotPluginWorkerRequest): Promise { + const worker = this.ensureWorker(); + return new Promise((resolve, reject) => { + this.pendingRequests.set(message.correlationId, { reject, resolve }); + worker.postMessage({ + message, + requestId: message.correlationId, + type: 'request', + }); + }); + } + + /** + * Terminates the worker thread and rejects in-flight requests for this runtime instance. + */ + async dispose(): Promise { + const worker = this.worker; + this.worker = undefined; + this.rejectPending(new Error('QQBot 插件 worker 已关闭')); + if (worker) { + await worker.terminate(); + } + } + + /** + * Lazily starts the worker thread with descriptor, installation, plugin key, and config snapshot workerData. + * @returns Active worker thread for the current descriptor runtime. + */ + private ensureWorker() { + if (this.worker) return this.worker; + + const worker = new Worker(resolveWorkerEntrypoint(), { + execArgv: resolveWorkerExecArgv(), + workerData: { + configSnapshot: this.options.configSnapshot, + descriptor: this.options.descriptor, + installationId: this.options.installationId, + pluginKey: this.options.pluginKey, + }, + }); + worker.on('message', (message: WorkerBridgeMessage) => { + void this.handleWorkerMessage(message); + }); + worker.on('error', (error) => { + this.worker = undefined; + this.rejectPending(error); + }); + worker.on('exit', (code) => { + this.worker = undefined; + if (code !== 0) { + this.rejectPending(new Error(`QQBot 插件 worker 异常退出:${code}`)); + } + }); + this.worker = worker; + return worker; + } + + /** + * Handles worker responses and host-call requests emitted by the child thread. + * @param message - Worker bridge message containing either a runtime response or a host capability request. + */ + private async handleWorkerMessage(message: WorkerBridgeMessage) { + if (message.type === 'response') { + this.settleWorkerResponse(message); + return; + } + + try { + const response = await this.hostBridge.handleHostCall( + this.options.descriptor, + { + args: message.args || {}, + method: message.method, + pluginKey: this.options.pluginKey, + }, + ); + if (response.ok === true) { + this.worker?.postMessage({ + ok: true, + requestId: message.requestId, + result: response.value, + type: 'hostResponse', + }); + return; + } + this.worker?.postMessage({ + error: { + message: + response.ok === false ? response.message : 'Host call failed', + name: 'QqbotPluginHostCallError', + }, + ok: false, + requestId: message.requestId, + type: 'hostResponse', + }); + } catch (error) { + this.worker?.postMessage({ + error: serializePluginWorkerResponseError(error), + ok: false, + requestId: message.requestId, + type: 'hostResponse', + }); + } + } + + /** + * Resolves or rejects one pending runtime request from a worker response message. + * @param message - Worker response carrying the original request id and serialized result or error. + */ + private settleWorkerResponse( + message: Extract, + ) { + const pending = this.pendingRequests.get(message.requestId); + if (!pending) return; + this.pendingRequests.delete(message.requestId); + if (message.ok) { + pending.resolve(message.result); + return; + } + pending.reject(new QqbotPluginWorkerResponseError(message.error || {})); + } + + /** + * Rejects all in-flight requests when the worker exits or is disposed. + * @param error - Runtime boundary error propagated to all pending callers. + */ + private rejectPending(error: Error) { + for (const pending of this.pendingRequests.values()) { + pending.reject(error); + } + this.pendingRequests.clear(); + } +} + +/** + * Resolves the generic worker thread entrypoint next to the compiled factory file. + * @returns Absolute worker entry file path with the current TypeScript or JavaScript extension. + */ +function resolveWorkerEntrypoint() { + const extension = __filename.endsWith('.ts') ? '.ts' : '.js'; + return join(__dirname, `plugin-worker.thread${extension}`); +} + +/** + * Resolves worker exec arguments needed when tests or local dev run TypeScript sources directly. + * @returns Node exec arguments that register ts-node only for `.ts` runtime files. + */ +function resolveWorkerExecArgv() { + if (!__filename.endsWith('.ts')) return []; + return ['-r', 'ts-node/register', '-r', 'tsconfig-paths/register']; +} diff --git a/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker.thread.ts b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker.thread.ts new file mode 100644 index 0000000..84489e6 --- /dev/null +++ b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker.thread.ts @@ -0,0 +1,640 @@ +import { createRequire } from 'node:module'; +import { readFileSync } from 'node:fs'; +import { dirname } from 'node:path'; +import { parentPort, workerData } from 'node:worker_threads'; +import { pathToFileURL } from 'node:url'; + +import type { + QqbotPluginPackageDescriptor, + QqbotPluginRuntimeConfigSnapshot, +} from '@/modules/qqbot/plugin-platform/infrastructure/integration/package/plugin-package.types'; +import type { QqbotPluginWorkerRequest } from './worker-runtime.types'; +import type { QqbotPluginWorkerRequestType } from './worker-runtime.types'; + +export type QqbotWorkerCreatePluginOptions = { + configSnapshot: QqbotPluginRuntimeConfigSnapshot; + descriptor: QqbotPluginPackageDescriptor; + host: Record; + installationId: string; +}; + +export type QqbotWorkerPluginInstance = { + activate?: () => Promise | unknown; + dispose?: () => Promise | unknown; + executeOperation?: ( + operationKey: string, + input: unknown, + ) => Promise | unknown; + executeTask?: ( + taskKey: string, + input: unknown, + context?: unknown, + ) => Promise | unknown; + handleEvent?: ( + eventKey: string, + event: unknown, + ) => Promise | unknown; + health?: () => Promise | unknown; + healthCheck?: () => Promise | unknown; + operations?: unknown[]; + tasks?: unknown[]; +}; + +type PluginEntryModule = { + createPlugin?: (input: { + host: Record; + manifest: QqbotPluginPackageDescriptor['manifest']; + normalizeError: (error: unknown, fallback?: string) => string; + now: () => Date; + runtime: { + configSnapshot: QqbotPluginRuntimeConfigSnapshot; + installationId: string; + }; + }) => Promise | QqbotWorkerPluginInstance; + default?: { + createPlugin?: PluginEntryModule['createPlugin']; + }; +}; + +type ParentMessage = + | { + message: QqbotPluginWorkerRequest; + requestId: string; + type: 'request'; + } + | { + error?: { message?: string; name?: string; stack?: string }; + ok: boolean; + requestId: string; + result?: unknown; + type: 'hostResponse'; + }; + +type PendingHostCall = { + reject: (reason?: unknown) => void; + resolve: (value: unknown) => void; +}; + +type HostArgumentMapper = (...args: unknown[]) => Record; + +const HOST_ARGUMENT_MAPPERS: Record = { + bindEventPlugin: (selfId, pluginKey) => ({ pluginKey, selfId }), + getBoundEventPluginKeys: (selfId) => ({ selfId }), + getConfig: (key) => ({ key }), + getConfigMany: (keys) => ({ keys }), + getDictByKey: (dictCode) => ({ dictCode }), + getDictItemsByKey: (dictCode) => ({ dictCode }), + readAssetFile: (path) => ({ path }), + readJsonFile: (path) => ({ path }), + relationTree: (input) => ({ input }), + renameFile: (from, to) => ({ from, to }), + requestBuffer: (options) => ({ options }), + requestJson: (options) => ({ options }), + sendText: (input) => ({ input }), + sleep: (ms) => ({ ms }), + unbindEventPlugin: (selfId, pluginKey) => ({ pluginKey, selfId }), + warn: (message) => ({ message }), + writeJsonFile: (path, data) => ({ data, path }), +}; + +const port = parentPort; +const pendingHostCalls = new Map(); +const requireEntryModule = createRequire(__filename); +let plugin: QqbotWorkerPluginInstance | undefined; + +const WORKER_REQUEST_HANDLERS: Record< + QqbotPluginWorkerRequestType, + (message: QqbotPluginWorkerRequest) => Promise | unknown +> = { + /** + * Loads a descriptor entry and stores the plugin instance for later requests. + * @param message - Load request carrying installation fallback metadata. + */ + load: (message) => loadPlugin(message), + /** + * Activates the loaded plugin instance. + */ + activate: async () => { + await requirePlugin().activate?.(); + return { ok: true }; + }, + /** + * Runs the loaded plugin health hook. + */ + health: () => healthPlugin(), + /** + * Executes one manifest operation. + * @param message - Operation request with operation key and input. + */ + executeOperation: (message) => executeOperation(message), + /** + * Executes one manifest task. + * @param message - Task request with task key and input. + */ + executeTask: (message) => executeTask(message), + /** + * Dispatches one manifest event. + * @param message - Event request with event key and payload. + */ + handleEvent: (message) => handleEvent(message), + /** + * Marks the worker inactive without disposing the plugin instance. + */ + deactivate: () => ({ ok: true }), + /** + * Disposes the loaded plugin instance and clears worker state. + */ + dispose: async () => { + await plugin?.dispose?.(); + plugin = undefined; + return { ok: true }; + }, +}; + +/** + * Dynamically imports a plugin package entry and creates its runtime instance. + * @param options - Descriptor, host facade, installation id, and config snapshot supplied by the worker factory. + * @returns Plugin instance returned by the package `createPlugin` entry. + */ +export async function createPluginFromDescriptor( + options: QqbotWorkerCreatePluginOptions, +): Promise { + const moduleUrl = pathToFileURL(options.descriptor.entryFile).href; + const entryModule = await importPluginEntryModule( + moduleUrl, + options.descriptor.entryFile, + ); + const createPlugin = + entryModule.createPlugin || entryModule.default?.createPlugin; + + if (typeof createPlugin !== 'function') { + throw new Error('Plugin entry must export createPlugin(options)'); + } + + return createPlugin({ + host: options.host, + manifest: options.descriptor.manifest, + normalizeError, + now: () => new Date(), + runtime: { + configSnapshot: options.configSnapshot, + installationId: options.installationId, + }, + }); +} + +if (port) { + port.on('message', (message: ParentMessage) => { + void handleParentMessage(message); + }); +} + +/** + * Handles a parent RPC or host-call response message inside the worker thread. + * @param message - Parent message carrying a worker request or the response for a pending host call. + */ +async function handleParentMessage(message: ParentMessage) { + if (message.type === 'hostResponse') { + settleHostResponse(message); + return; + } + + try { + const result = await handleWorkerRequest(message.message); + port?.postMessage({ + ok: true, + requestId: message.requestId, + result, + type: 'response', + }); + } catch (error) { + port?.postMessage({ + error: serializeError(error), + ok: false, + requestId: message.requestId, + type: 'response', + }); + } +} + +/** + * Dispatches a lifecycle or execution request to the loaded plugin instance. + * @param message - Worker request produced by QqbotPluginWorkerRuntime. + * @returns Plugin lifecycle, operation, task, or event result. + */ +async function handleWorkerRequest(message: QqbotPluginWorkerRequest) { + return WORKER_REQUEST_HANDLERS[message.type](message); +} + +/** + * Creates and stores the plugin instance for subsequent worker requests. + * @param message - Load request carrying the installation id fallback used by the worker runtime. + * @returns Safe load summary for the parent runtime. + */ +async function loadPlugin(message: QqbotPluginWorkerRequest) { + const descriptor = getWorkerDescriptor(); + plugin = await createPluginFromDescriptor({ + configSnapshot: getWorkerConfigSnapshot(), + descriptor, + host: createHostFacade(), + installationId: getWorkerInstallationId(message), + }); + + return { + ok: true, + pluginKey: descriptor.pluginKey, + }; +} + +/** + * Runs the plugin health hook using either modern `health` or legacy `healthCheck`. + * @returns Health payload produced by the plugin, or a generic healthy status. + */ +async function healthPlugin() { + const loadedPlugin = requirePlugin(); + if (loadedPlugin.health) return loadedPlugin.health(); + if (loadedPlugin.healthCheck) return loadedPlugin.healthCheck(); + return { ok: true, status: 'healthy' }; +} + +/** + * Executes one plugin operation through the generic instance contract. + * @param message - Operation request with manifest operation key and user input payload. + * @returns Operation output produced by the plugin. + */ +async function executeOperation(message: QqbotPluginWorkerRequest) { + const loadedPlugin = requirePlugin(); + const operationKey = requireRequestKey( + message.operationKey || message.operationId, + 'QQBot 插件能力缺少 operationKey', + ); + + if (loadedPlugin.executeOperation) { + return loadedPlugin.executeOperation(operationKey, message.input); + } + + const operation = findRuntimeCallable(loadedPlugin.operations, operationKey); + if (operation) { + return operation.execute(message.input || {}); + } + + throw new Error(`QQBot 插件能力不存在:${operationKey}`); +} + +/** + * Executes one plugin task through the generic instance contract. + * @param message - Task request with task key, handler name, trigger type, and input payload. + * @returns Task output produced by the plugin. + */ +async function executeTask(message: QqbotPluginWorkerRequest) { + const loadedPlugin = requirePlugin(); + const taskKey = requireRequestKey( + message.taskKey || message.taskHandlerName, + 'QQBot 插件定时任务缺少 taskKey', + ); + const context = { + taskHandlerName: message.taskHandlerName, + taskId: message.taskId, + triggerType: message.triggerType, + }; + + if (loadedPlugin.executeTask) { + return loadedPlugin.executeTask(taskKey, message.input, context); + } + + const task = findRuntimeCallable( + loadedPlugin.tasks, + taskKey, + message.taskHandlerName, + ); + if (task) { + return task.execute(message.input || {}, context); + } + + throw new Error(`QQBot 插件定时任务不存在:${taskKey}`); +} + +/** + * Dispatches one plugin event through the generic instance contract. + * @param message - Event request with event key and normalized event payload. + * @returns Event handling result produced by the plugin. + */ +async function handleEvent(message: QqbotPluginWorkerRequest) { + const loadedPlugin = requirePlugin(); + const eventKey = requireRequestKey( + message.eventKey, + 'QQBot 插件事件缺少 eventKey', + ); + + if (!loadedPlugin.handleEvent) return false; + return loadedPlugin.handleEvent(eventKey, message.event); +} + +/** + * Creates a host facade whose methods proxy generic host capabilities to the parent process. + * @returns Record of host capability functions passed to plugin package entries. + */ +function createHostFacade(): Record { + const host: Record = {}; + for (const method of Object.keys(HOST_ARGUMENT_MAPPERS)) { + host[method] = (...args: unknown[]) => + callHost(method, HOST_ARGUMENT_MAPPERS[method](...args)); + } + + return new Proxy(host, { + /** + * Resolves unknown host methods to generic RPC calls while keeping promise introspection stable. + * @param target - Known host capability map. + * @param property - Host method name requested by plugin package code. + * @returns Existing method, generic method proxy, or `undefined` for non-method symbols. + */ + get(target, property) { + if (typeof property !== 'string') return undefined; + if (property in target) return target[property]; + if (property === 'then') return undefined; + return (...args: unknown[]) => + callHost(property, normalizeHostArgs(args)); + }, + }); +} + +/** + * Sends one host capability request to the parent process and waits for its response. + * @param method - Generic host capability name. + * @param args - Structured argument payload for the host bridge. + * @returns Host bridge result value. + */ +function callHost( + method: string, + args: Record, +): Promise { + if (!port) { + return Promise.reject( + new Error('QQBot plugin worker host port unavailable'), + ); + } + + const requestId = `host-${Date.now()}-${Math.random().toString(16).slice(2)}`; + return new Promise((resolve, reject) => { + pendingHostCalls.set(requestId, { + reject, + /** + * Resolves a host response while preserving the generic result type requested by plugin code. + * @param value - Raw host bridge response value sent by the parent process. + */ + resolve: (value) => resolve(value as TResult), + }); + port.postMessage({ + args, + method, + requestId, + type: 'hostCall', + }); + }); +} + +/** + * Resolves or rejects a pending host call from a parent host response message. + * @param message - Host response carrying the original host-call request id. + */ +function settleHostResponse( + message: Extract, +) { + const pending = pendingHostCalls.get(message.requestId); + if (!pending) return; + pendingHostCalls.delete(message.requestId); + + if (message.ok) { + pending.resolve(message.result); + return; + } + + pending.reject(deserializeError(message.error)); +} + +/** + * Reads the workerData descriptor and rejects malformed worker startup state. + * @returns Package descriptor supplied by QqbotPluginWorkerThreadDriver. + */ +function getWorkerDescriptor(): QqbotPluginPackageDescriptor { + const descriptor = workerData?.descriptor; + if (!descriptor || typeof descriptor !== 'object') { + throw new Error('QQBot 插件 worker 缺少 descriptor'); + } + return descriptor as QqbotPluginPackageDescriptor; +} + +/** + * Reads the workerData config snapshot as a simple string map. + * @returns Runtime config snapshot supplied by QqbotPluginWorkerThreadDriver. + */ +function getWorkerConfigSnapshot(): QqbotPluginRuntimeConfigSnapshot { + const snapshot = workerData?.configSnapshot; + return snapshot && typeof snapshot === 'object' + ? (snapshot as QqbotPluginRuntimeConfigSnapshot) + : {}; +} + +/** + * Reads the installation id from workerData, falling back to the load request for recovery calls. + * @param message - Load request that may carry the installation id from QqbotPluginWorkerRuntime. + * @returns Installation id passed to the plugin runtime contract. + */ +function getWorkerInstallationId(message: QqbotPluginWorkerRequest) { + const installationId = workerData?.installationId || message.installationId; + return typeof installationId === 'string' ? installationId : ''; +} + +/** + * Returns the currently loaded plugin instance or throws a worker-safe error. + * @returns Plugin instance created during the load request. + */ +function requirePlugin() { + if (!plugin) { + throw new Error('QQBot 插件运行时未加载'); + } + return plugin; +} + +/** + * Reads a required request key from a worker request. + * @param value - Candidate request key from operation, task, or event payload. + * @param message - Error message used when the key is absent. + * @returns Non-empty request key. + */ +function requireRequestKey(value: unknown, message: string) { + if (typeof value === 'string' && value.trim()) return value.trim(); + throw new Error(message); +} + +/** + * Finds an executable operation or task from a manifest-shaped runtime summary. + * @param items - Runtime operation or task list returned by plugin entry code. + * @param key - Operation or task key requested by the platform. + * @param handlerName - Optional handler name fallback for task execution. + * @returns Callable runtime item when found. + */ +function findRuntimeCallable( + items: unknown[] | undefined, + key: string, + handlerName?: string, +) { + return items + ?.filter(isRuntimeCallable) + .find( + (item) => + item.key === key || + item.handlerName === key || + (!!handlerName && item.handlerName === handlerName), + ); +} + +/** + * Checks whether a runtime summary item exposes an executable function. + * @param item - Unknown operation or task summary returned by plugin code. + * @returns `true` when the item can execute worker input. + */ +function isRuntimeCallable(item: unknown): item is { + execute: (input: unknown, context?: unknown) => Promise | unknown; + handlerName?: string; + key?: string; +} { + return ( + !!item && + typeof item === 'object' && + typeof (item as { execute?: unknown }).execute === 'function' + ); +} + +/** + * Normalizes unknown host arguments for unlisted host methods. + * @param args - Positional arguments supplied by plugin package code. + * @returns Record payload accepted by the parent host bridge. + */ +function normalizeHostArgs(args: unknown[]): Record { + if (args.length === 1 && isRecord(args[0])) return args[0]; + return { args }; +} + +/** + * Checks whether a value can be passed as a structured host-call argument record. + * @param value - Candidate host-call argument value. + * @returns `true` when the value is a non-null object and not an array. + */ +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value); +} + +/** + * Imports a plugin entry with a URL-based dynamic import and falls back to CommonJS loading in Jest/dev hooks. + * @param moduleUrl - File URL produced from the descriptor entry path. + * @param entryFile - Absolute descriptor entry file used only when the dynamic import host cannot resolve file URLs. + * @returns Plugin entry module namespace. + */ +async function importPluginEntryModule(moduleUrl: string, entryFile: string) { + try { + return await dynamicImportPluginEntry(moduleUrl); + } catch (error) { + if (!shouldFallbackToRequire(error, moduleUrl)) { + throw error; + } + if (entryFile.endsWith('.js') || entryFile.endsWith('.cjs')) { + return loadCommonJsEntryModule(entryFile); + } + return requireEntryModule(entryFile) as PluginEntryModule; + } +} + +/** + * Uses native dynamic import so Jest CommonJS transforms do not rewrite file URL imports. + * @param moduleUrl - File URL pointing at the descriptor entry module. + * @returns Plugin entry module namespace loaded by Node. + */ +function dynamicImportPluginEntry( + moduleUrl: string, +): Promise { + const importer = new Function('moduleUrl', 'return import(moduleUrl)') as ( + value: string, + ) => Promise; + return importer(moduleUrl); +} + +/** + * Detects Jest/ts-node resolution failures for URL imports while preserving real plugin errors. + * @param error - Dynamic import error thrown by the runtime host. + * @param moduleUrl - File URL attempted by the dynamic import. + * @returns Whether CommonJS fallback is appropriate for the same entry file. + */ +function shouldFallbackToRequire(error: unknown, moduleUrl: string) { + if (!moduleUrl.startsWith('file:')) return false; + const message = error instanceof Error ? error.message : `${error}`; + return ( + message.includes('Cannot find module') || + message.includes('dynamic import callback') || + message.includes('Unknown file extension') + ); +} + +/** + * Loads a CommonJS plugin entry directly from disk for Jest environments that cannot execute file URL imports. + * @param entryFile - Absolute descriptor entry file compiled as CommonJS. + * @returns Plugin entry module exports. + */ +function loadCommonJsEntryModule(entryFile: string): PluginEntryModule { + const nodeModule = requireEntryModule( + 'node:module', + ) as typeof import('node:module'); + const moduleConstructor = nodeModule.Module as typeof nodeModule.Module & { + _nodeModulePaths: (from: string) => string[]; + }; + const entryModule = new moduleConstructor(entryFile); + entryModule.filename = entryFile; + entryModule.paths = moduleConstructor._nodeModulePaths(dirname(entryFile)); + ( + entryModule as typeof entryModule & { + _compile: (content: string, filename: string) => void; + } + )._compile(readFileSync(entryFile, 'utf8'), entryFile); + return entryModule.exports as PluginEntryModule; +} + +/** + * Normalizes plugin errors into stable text for package entry adapters. + * @param error - Error or arbitrary thrown value from package code. + * @param fallback - Message used when the thrown value has no useful text. + * @returns Stable error message for plugin-visible health or execution output. + */ +function normalizeError(error: unknown, fallback = '插件执行失败') { + if (error instanceof Error && error.message) return error.message; + const message = `${error || ''}`.trim(); + return message || fallback; +} + +/** + * Serializes an error for parent-process worker responses. + * @param error - Error or arbitrary thrown value from worker or plugin code. + * @returns Worker-safe error payload with message, name, and stack when available. + */ +function serializeError(error: unknown) { + return { + message: error instanceof Error ? error.message : `${error}`, + name: error instanceof Error ? error.name : 'Error', + stack: error instanceof Error ? error.stack : undefined, + }; +} + +/** + * Deserializes a parent host-call error into an Error instance. + * @param error - Serialized error payload from the parent host bridge. + * @returns Error instance rejected to plugin host facade callers. + */ +function deserializeError(error?: { + message?: string; + name?: string; + stack?: string; +}) { + const output = new Error(error?.message || '插件 Host 调用失败'); + if (error?.name) output.name = error.name; + if (error?.stack) output.stack = error.stack; + return output; +} diff --git a/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/worker-runtime.types.ts b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/worker-runtime.types.ts index 9f8f333..22d64f3 100644 --- a/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/worker-runtime.types.ts +++ b/src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/worker-runtime.types.ts @@ -59,6 +59,7 @@ export type QqbotPluginWorkerRequestQueue = { }; export type QqbotPluginWorkerRuntimeOptions = { + configSnapshot?: QqbotPluginRuntimeConfigSnapshot; defaultTimeoutMs: number; descriptor?: QqbotPluginPackageDescriptor; installationId: string; diff --git a/test/modules/qqbot/plugin-platform/plugin-lifecycle-runtime.spec.ts b/test/modules/qqbot/plugin-platform/plugin-lifecycle-runtime.spec.ts index 95b3ff5..53f1231 100644 --- a/test/modules/qqbot/plugin-platform/plugin-lifecycle-runtime.spec.ts +++ b/test/modules/qqbot/plugin-platform/plugin-lifecycle-runtime.spec.ts @@ -8,7 +8,10 @@ import { QqbotEventPluginRegistryService } from '../../../../src/modules/qqbot/p import { QqbotPluginRegistryService } from '../../../../src/modules/qqbot/plugin-platform/application/registry/qqbot-plugin-registry.service'; import { QqbotPluginPlatformModule } from '../../../../src/modules/qqbot/plugin-platform/plugin-platform.module'; import { QqbotBuiltinPluginPackageLoaderService } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/package/builtin-plugin-package-loader.service'; +import { QqbotPluginPackageSourceService } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/package/plugin-package-source.service'; +import { QqbotPluginHostBridgeService } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-host-bridge.service'; import { QqbotBuiltinPluginWorkerRuntimeFactoryService } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/runtime'; +import { QqbotPluginWorkerRuntimeFactoryService } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/runtime'; const repoRoot = join(__dirname, '../../../..'); @@ -62,6 +65,35 @@ describe('QQBot plugin platform lifecycle runtime contract', () => { expect(dependencies?.[1]).toBe(ConfigService); }); + it('keeps the generic worker runtime factory dependency visible to Nest DI', () => { + const dependencies = Reflect.getMetadata( + 'design:paramtypes', + QqbotPluginWorkerRuntimeFactoryService, + ); + + expect(dependencies?.[0]).toBe(ConfigService); + expect(dependencies?.[1]).toBe(QqbotPluginPackageSourceService); + expect(dependencies?.[2]).toBe(QqbotPluginHostBridgeService); + }); + + it('keeps the generic worker runtime descriptor-based and plugin-agnostic', () => { + const source = [ + readSource( + 'src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker-runtime.factory.ts', + ), + readSource( + 'src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker.thread.ts', + ), + ].join('\n'); + + expect(source).toContain('node:worker_threads'); + expect(source).toContain('descriptor'); + expect(source).not.toContain('QqbotBuiltinPluginPackageLoaderService'); + expect(source).not.toMatch( + /@\/modules\/qqbot\/plugins|src\/modules\/qqbot\/plugins|createBangDreamPlugin|createFf14MarketPlugin|createFflogsPlugin|createRepeaterPlugin|pluginKey\s*===|case\s+['"]/, + ); + }); + it('uses a real worker-thread boundary for built-in plugin runtimes', () => { const source = readSource( 'src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/builtin-plugin-worker-runtime.factory.ts', diff --git a/test/modules/qqbot/plugin-platform/plugin-worker-entry.spec.ts b/test/modules/qqbot/plugin-platform/plugin-worker-entry.spec.ts new file mode 100644 index 0000000..c772dc1 --- /dev/null +++ b/test/modules/qqbot/plugin-platform/plugin-worker-entry.spec.ts @@ -0,0 +1,114 @@ +import { + mkdtempSync, + mkdirSync, + readFileSync, + rmSync, + writeFileSync, +} from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; + +import type { QqbotPluginManifest } from '@/modules/qqbot/plugin-platform/domain/manifest'; +import { createPluginFromDescriptor } from '@/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker.thread'; + +describe('generic plugin worker entry', () => { + let tempRoot: string; + + beforeEach(() => { + tempRoot = mkdtempSync(join(tmpdir(), 'qqbot-plugin-worker-')); + }); + + afterEach(() => { + rmSync(tempRoot, { force: true, recursive: true }); + }); + + it('loads a package entry from the descriptor instead of a plugin key switch', async () => { + const packageRoot = join(tempRoot, 'sample'); + mkdirSync(join(packageRoot, 'src'), { recursive: true }); + const entryFile = join(packageRoot, 'src', 'index.js'); + writeFileSync( + entryFile, + ` + exports.createPlugin = (options) => ({ + activate: async () => undefined, + health: async () => ({ + configSnapshot: options.runtime.configSnapshot, + installationId: options.runtime.installationId, + ok: true, + pluginKey: options.manifest.key, + }), + operations: [], + }); + `, + 'utf8', + ); + + const plugin = await createPluginFromDescriptor({ + configSnapshot: { SAMPLE_TOKEN: 'local-token' }, + descriptor: { + entry: 'src/index.js', + entryFile, + manifest: createManifest('sample', 'src/index.js'), + packageRoot, + pluginKey: 'sample', + }, + host: {}, + installationId: 'install-1', + }); + + await expect(plugin.health?.()).resolves.toEqual({ + configSnapshot: { SAMPLE_TOKEN: 'local-token' }, + installationId: 'install-1', + ok: true, + pluginKey: 'sample', + }); + }); + + it('keeps the generic worker entry free of concrete plugin branches', () => { + const source = readFileSync( + join( + process.cwd(), + 'src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/plugin-worker.thread.ts', + ), + 'utf8', + ); + + expect(source).toContain('node:worker_threads'); + expect(source).toContain('descriptor'); + expect(source).not.toMatch( + /createBangDreamPlugin|createFf14MarketPlugin|createFflogsPlugin|createRepeaterPlugin|pluginKey\s*===|case\s+['"]/, + ); + }); +}); + +/** + * Builds a minimal parsed-manifest-shaped descriptor manifest for direct worker entry tests. + * @param key - Plugin package key used by the temp descriptor and health assertion. + * @param entry - Package-relative entry file that the descriptor resolves to the temp module. + * @returns Manifest payload accepted by descriptor-based worker creation. + */ +function createManifest(key: string, entry: string): QqbotPluginManifest { + return { + assets: [], + configSchema: {}, + entry, + events: [], + key, + legacyAliases: [], + migrations: [], + minApiSdkVersion: '1.0.0', + name: 'Sample', + operations: [], + permissions: [], + pluginKey: key, + runtime: { + configKeys: ['SAMPLE_TOKEN'], + maxConcurrency: 1, + memoryMb: 128, + timeoutMs: 5000, + workerType: 'thread', + }, + tasks: [], + version: '1.0.0', + }; +}