import type { ConfigService } from '@nestjs/config'; import { QqbotPluginRuntimeError, QqbotPluginWorkerRuntime, QqbotPluginWorkerExpiredRequestError, QqbotPluginWorkerResponseError, } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/worker-runtime'; import { createQqbotBullmqWorkerQueueOptions, resolveQqbotPluginQueueConnection, } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/bullmq-plugin-worker-request.queue'; import type { QqbotPluginWorkerDriver, QqbotPluginWorkerRequest, QqbotPluginWorkerRequestQueue, QqbotPluginWorkerRuntimeOptions, } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/runtime/worker-runtime.types'; import { createQqbotPluginSdk } from '../../../../src/modules/qqbot/plugin-platform/infrastructure/integration/sdk'; class RecordingDriver implements QqbotPluginWorkerDriver { disposed = false; readonly requests: QqbotPluginWorkerRequest[] = []; readonly responses = new Map(); constructor(private readonly rejectTypes = new Map()) {} /** * 执行 QQBot 插件平台流程。 */ async dispose() { this.disposed = true; } /** * 执行 QQBot 插件平台流程。 * @param message - message 输入;使用 `type` 字段生成结果。 * @returns 异步完成后的 QQBot 插件平台结果。 */ async request(message: QqbotPluginWorkerRequest): Promise { this.requests.push(message); const rejection = this.rejectTypes.get(message.type); if (rejection) throw rejection; if (this.responses.has(message.type)) { return this.responses.get(message.type); } return { ok: true, type: message.type, }; } } class RecordingRequestQueue implements QqbotPluginWorkerRequestQueue { private previous = Promise.resolve(); constructor(private readonly driver: QqbotPluginWorkerDriver) {} /** * 执行 QQBot 插件平台流程。 */ async close() { await this.driver.dispose(); } /** * 执行 QQBot 插件平台流程。 * @param message - message 输入;驱动 `catch()` 的 插件平台步骤。 * @returns 异步完成后的 QQBot 插件平台结果。 */ request(message: QqbotPluginWorkerRequest): Promise { const run = this.previous .catch(() => undefined) .then(() => this.driver.request(message)); this.previous = run.then( () => undefined, () => undefined, ); return run; } /** * 重置业务数据。 */ async reset() { await this.driver.dispose(); } } class GenerationAwareRequestQueue implements QqbotPluginWorkerRequestQueue { private generation = 0; private previous = Promise.resolve(); constructor(private readonly driver: QqbotPluginWorkerDriver) {} /** * 执行 QQBot 插件平台流程。 */ async close() { await this.driver.dispose(); } /** * 执行 QQBot 插件平台流程。 * @param message - message 输入;驱动 `driver.request()` 的 插件平台步骤。 * @returns 异步完成后的 QQBot 插件平台结果。 */ request(message: QqbotPluginWorkerRequest): Promise { const generation = this.generation; const run = this.previous .catch(() => undefined) .then(async () => { if (generation !== this.generation) { const error = new Error('stale worker queue request'); error.name = 'QqbotPluginWorkerStaleRequestError'; throw error; } try { return await this.driver.request(message); } catch (error) { if (!(error instanceof QqbotPluginWorkerResponseError)) { this.generation += 1; } throw error; } }); this.previous = run.then( () => undefined, () => undefined, ); return run; } /** * 重置业务数据。 */ async reset() { this.generation += 1; await this.driver.dispose(); } } class TimeoutAwareRecordingRequestQueue extends RecordingRequestQueue { readonly handlesRequestTimeout = true; readonly queueWaitTimeoutMs = 50; } class ExpiringOnceRequestQueue implements QqbotPluginWorkerRequestQueue { constructor(private readonly driver: QqbotPluginWorkerDriver) {} /** * 执行 QQBot 插件平台流程。 */ async close() { await this.driver.dispose(); } /** * 执行 QQBot 插件平台流程。 * @param message - message 输入;使用 `operationId` 字段生成结果。 * @returns 异步完成后的 QQBot 插件平台结果。 */ async request(message: QqbotPluginWorkerRequest): Promise { if (message.operationId === 'op-expire') { await this.reset(); throw new QqbotPluginWorkerExpiredRequestError( 'worker-request-execution-timeout', ); } return this.driver.request(message); } /** * 重置业务数据。 */ async reset() { await this.driver.dispose(); } } const createRuntime = (driver = new RecordingDriver()) => { const runtime = new QqbotPluginWorkerRuntime( new RecordingRequestQueue(driver), { defaultTimeoutMs: 50, installationId: 'install-1', pluginKey: 'demo-plugin', }, ); return { driver, runtime }; }; describe('QQBot plugin worker runtime', () => { it('accepts a package descriptor in worker runtime options', () => { const options: QqbotPluginWorkerRuntimeOptions = { defaultTimeoutMs: 5000, descriptor: { entry: 'src/index.ts', entryFile: 'D:/repo/src/modules/qqbot/plugins/sample/src/index.ts', manifest: { key: 'sample', pluginKey: 'sample', name: 'Sample', version: '1.0.0', minApiSdkVersion: '1.0.0', entry: 'src/index.ts', runtime: { configKeys: ['SAMPLE_TOKEN'], maxConcurrency: 1, memoryMb: 128, timeoutMs: 5000, workerType: 'thread', }, operations: [], events: [], tasks: [], assets: [], configSchema: {}, legacyAliases: [], migrations: [], permissions: [], }, packageRoot: 'D:/repo/src/modules/qqbot/plugins/sample', pluginKey: 'sample', }, installationId: 'install-1', pluginKey: 'sample', }; expect(options.descriptor.manifest.runtime.configKeys).toEqual([ 'SAMPLE_TOKEN', ]); expect(options.descriptor.pluginKey).toBe('sample'); }); it('serializes concurrent worker execution requests for one plugin runtime', async () => { const releaseFirst = createDeferred(); const releaseSecond = createDeferred(); let activeRequests = 0; let maxActiveRequests = 0; const requestOrder: string[] = []; const driver: QqbotPluginWorkerDriver = { dispose: jest.fn(async () => undefined), request: jest.fn(async (message) => { activeRequests += 1; maxActiveRequests = Math.max(maxActiveRequests, activeRequests); requestOrder.push(message.operationId || message.type); if (message.operationId === 'op-1') { await releaseFirst.promise; } if (message.operationId === 'op-2') { await releaseSecond.promise; } activeRequests -= 1; return { ok: true, operationId: message.operationId, }; }), }; const runtime = new QqbotPluginWorkerRuntime( new RecordingRequestQueue(driver), { defaultTimeoutMs: 100, installationId: 'install-serial', pluginKey: 'demo-plugin', }, ); const first = runtime.executeOperation({ input: { text: 'first' }, operationId: 'op-1', operationKey: 'demo-plugin.echo', }); const second = runtime.executeOperation({ input: { text: 'second' }, operationId: 'op-2', operationKey: 'demo-plugin.echo', }); try { await waitUntil(() => requestOrder.includes('op-1')); expect(requestOrder).toEqual(['op-1']); releaseFirst.resolve(); await expect(first).resolves.toMatchObject({ operationId: 'op-1' }); await waitUntil(() => requestOrder.includes('op-2')); expect(maxActiveRequests).toBe(1); releaseSecond.resolve(); await expect(second).resolves.toMatchObject({ operationId: 'op-2' }); } finally { releaseFirst.resolve(); releaseSecond.resolve(); await Promise.allSettled([first, second]); } }); it('does not count queue wait time against a queued operation execution timeout', async () => { const requestOrder: string[] = []; const driver: QqbotPluginWorkerDriver = { dispose: jest.fn(async () => undefined), request: jest.fn(async (message) => { requestOrder.push(message.operationId || message.type); await new Promise((resolve) => setTimeout(resolve, message.operationId === 'op-1' ? 20 : 15), ); return { ok: true, operationId: message.operationId, }; }), }; const runtime = new QqbotPluginWorkerRuntime( new TimeoutAwareRecordingRequestQueue(driver), { defaultTimeoutMs: 30, installationId: 'install-queue-wait', pluginKey: 'demo-plugin', }, ); const first = runtime.executeOperation({ input: { text: 'first' }, operationId: 'op-1', operationKey: 'demo-plugin.echo', }); const second = runtime.executeOperation({ input: { text: 'second' }, operationId: 'op-2', operationKey: 'demo-plugin.echo', }); await expect(Promise.all([first, second])).resolves.toEqual([ { ok: true, operationId: 'op-1' }, { ok: true, operationId: 'op-2' }, ]); expect(requestOrder).toEqual(['op-1', 'op-2']); expect(driver.dispose).not.toHaveBeenCalled(); }); it('recovers a failed worker by reloading and activating it before the next request', async () => { const requestTypes: string[] = []; const driver: QqbotPluginWorkerDriver = { dispose: jest.fn(async () => undefined), request: jest.fn(async (message) => { requestTypes.push(message.type); if (message.operationId === 'op-crash') { throw new Error('worker crashed'); } return { ok: true, operationId: message.operationId, type: message.type, }; }), }; const runtime = new QqbotPluginWorkerRuntime( new RecordingRequestQueue(driver), { defaultTimeoutMs: 50, installationId: 'install-recover', pluginKey: 'demo-plugin', }, ); const manifest = { entry: 'src/index.ts', pluginKey: 'demo-plugin', version: '0.1.0', }; await runtime.load(manifest); await runtime.activate(); await expect( runtime.executeOperation({ input: { text: 'boom' }, operationId: 'op-crash', operationKey: 'demo-plugin.echo', }), ).rejects.toMatchObject({ code: 'PLUGIN_WORKER_CRASH', }); expect(runtime.status).toBe('failed'); await expect( runtime.executeOperation({ input: { text: 'after crash' }, operationId: 'op-after-crash', operationKey: 'demo-plugin.echo', }), ).resolves.toMatchObject({ operationId: 'op-after-crash' }); expect(driver.dispose).toHaveBeenCalled(); expect(requestTypes).toEqual([ 'load', 'activate', 'executeOperation', 'load', 'activate', 'executeOperation', ]); expect(runtime.status).toBe('active'); }); it('recovers before retrying requests that were already queued when a worker crashed', async () => { const requestTypes: string[] = []; const driver: QqbotPluginWorkerDriver = { dispose: jest.fn(async () => undefined), request: jest.fn(async (message) => { requestTypes.push(message.operationId || message.type); if (message.operationId === 'op-crash') { throw new Error('worker crashed'); } return { ok: true, operationId: message.operationId, type: message.type, }; }), }; const runtime = new QqbotPluginWorkerRuntime( new GenerationAwareRequestQueue(driver), { defaultTimeoutMs: 50, installationId: 'install-stale-retry', pluginKey: 'demo-plugin', }, ); await runtime.load({ entry: 'src/index.ts', pluginKey: 'demo-plugin', version: '0.1.0', }); await runtime.activate(); const first = runtime .executeOperation({ input: { text: 'boom' }, operationId: 'op-crash', operationKey: 'demo-plugin.echo', }) .catch((error) => error); const second = runtime.executeOperation({ input: { text: 'after crash' }, operationId: 'op-after-crash', operationKey: 'demo-plugin.echo', }); await expect(first).resolves.toMatchObject({ code: 'PLUGIN_WORKER_CRASH', }); await expect(second).resolves.toMatchObject({ operationId: 'op-after-crash', }); expect(requestTypes).toEqual([ 'load', 'activate', 'op-crash', 'load', 'activate', 'op-after-crash', ]); expect(runtime.status).toBe('active'); }); it('recovers after the queue expires and disposes a slow worker request', async () => { const requestTypes: string[] = []; const driver: QqbotPluginWorkerDriver = { dispose: jest.fn(async () => undefined), request: jest.fn(async (message) => { requestTypes.push(message.operationId || message.type); return { ok: true, operationId: message.operationId, type: message.type, }; }), }; const runtime = new QqbotPluginWorkerRuntime( new ExpiringOnceRequestQueue(driver), { defaultTimeoutMs: 50, installationId: 'install-expired-recover', pluginKey: 'demo-plugin', }, ); await runtime.load({ entry: 'src/index.ts', pluginKey: 'demo-plugin', version: '0.1.0', }); await runtime.activate(); await expect( runtime.executeOperation({ input: { text: 'slow' }, operationId: 'op-expire', operationKey: 'demo-plugin.echo', }), ).rejects.toMatchObject({ code: 'PLUGIN_WORKER_TIMEOUT', }); expect(runtime.status).toBe('failed'); await expect( runtime.executeOperation({ input: { text: 'after expired' }, operationId: 'op-after-expire', operationKey: 'demo-plugin.echo', }), ).resolves.toMatchObject({ operationId: 'op-after-expire', }); expect(requestTypes).toEqual([ 'load', 'activate', 'load', 'activate', 'op-after-expire', ]); expect(runtime.status).toBe('active'); }); it('sends lifecycle and execution RPC messages with correlation IDs and safe input summaries', async () => { const { driver, runtime } = createRuntime(); driver.responses.set('executeOperation', { replyText: 'ok' }); driver.responses.set('handleEvent', { handled: true }); await runtime.load({ entry: 'src/index.ts', pluginKey: 'demo-plugin', version: '0.1.0', }); await runtime.activate(); await expect( runtime.executeOperation({ input: { authMarker: 'sample-value', text: 'hello', }, operationId: 'op-1', operationKey: 'demo-plugin.echo', timeoutMs: 30, }), ).resolves.toEqual({ replyText: 'ok' }); await expect( runtime.handleEvent({ event: { message: 'hello', rawEvent: { traceMarker: 'sample-event', }, }, eventKey: 'demo-plugin.message', }), ).resolves.toEqual({ handled: true }); await runtime.health(); await runtime.deactivate(); await runtime.dispose(); expect(driver.requests.map((request) => request.type)).toEqual([ 'load', 'activate', 'executeOperation', 'handleEvent', 'health', 'deactivate', 'dispose', ]); expect(driver.requests.every((request) => request.correlationId)).toBe( true, ); expect(driver.requests[2]).toMatchObject({ operationId: 'op-1', operationKey: 'demo-plugin.echo', safeInputSummary: { fieldCount: 2, keys: ['authMarker', 'text'], }, timeoutMs: 30, type: 'executeOperation', }); expect(JSON.stringify(driver.requests[2].safeInputSummary)).not.toContain( 'sample-value', ); expect(JSON.stringify(driver.requests[3].safeInputSummary)).not.toContain( 'sample-event', ); expect(driver.disposed).toBe(true); }); it('sends executeTask RPC with safe input summary and timeout', async () => { const { driver, runtime } = createRuntime(); driver.responses.set('executeTask', { syncedKeys: ['songs'] }); await expect( runtime.executeTask({ input: { force: true, fullPayload: 'secret' }, taskHandlerName: 'syncBestdoriMainData', taskId: 'task-1', taskKey: 'bangdream.bestdori.sync-main-data', timeoutMs: 120000, triggerType: 'manual', }), ).resolves.toEqual({ syncedKeys: ['songs'] }); expect(driver.requests[0]).toMatchObject({ safeInputSummary: { fieldCount: 2, keys: ['force', 'fullPayload'] }, taskHandlerName: 'syncBestdoriMainData', taskId: 'task-1', taskKey: 'bangdream.bestdori.sync-main-data', timeoutMs: 120000, triggerType: 'manual', type: 'executeTask', }); expect(JSON.stringify(driver.requests[0].safeInputSummary)).not.toContain( 'secret', ); }); it('isolates worker crashes as plugin runtime events without throwing raw errors', async () => { const { runtime } = createRuntime( new RecordingDriver( new Map([['executeOperation', new Error('worker crashed')]]), ), ); await expect( runtime.executeOperation({ input: { text: 'boom' }, operationId: 'op-crash', operationKey: 'demo-plugin.echo', }), ).rejects.toMatchObject({ code: 'PLUGIN_WORKER_CRASH', pluginKey: 'demo-plugin', }); expect(runtime.status).toBe('failed'); expect(runtime.listRuntimeEvents()).toEqual([ expect.objectContaining({ eventType: 'worker-crash', level: 'error', safeSummary: expect.objectContaining({ message: 'worker crashed', operationId: 'op-crash', }), }), ]); }); it('keeps plugin response errors from poisoning the worker runtime', async () => { const driver = new RecordingDriver( new Map([ [ 'executeOperation', new QqbotPluginWorkerResponseError({ message: '业务参数错误', name: 'PluginCommandError', }), ], ]), ); const { runtime } = createRuntime(driver); await runtime.load({ entry: 'src/index.ts', pluginKey: 'demo-plugin', version: '0.1.0', }); await runtime.activate(); await expect( runtime.executeOperation({ input: { text: 'bad input' }, operationId: 'op-plugin-error', operationKey: 'demo-plugin.echo', }), ).rejects.toMatchObject({ message: '业务参数错误', name: 'PluginCommandError', }); expect(runtime.status).toBe('active'); expect(driver.disposed).toBe(false); expect(runtime.listRuntimeEvents()).toEqual([]); }); it('times out slow worker calls and disposes the driver boundary', async () => { const driver: QqbotPluginWorkerDriver = { dispose: jest.fn(async () => undefined), request: jest.fn( () => new Promise((resolve) => setTimeout(resolve, 100)), ), }; const runtime = new QqbotPluginWorkerRuntime( new RecordingRequestQueue(driver), { defaultTimeoutMs: 5, installationId: 'install-timeout', pluginKey: 'demo-plugin', }, ); await expect( runtime.executeOperation({ input: { text: 'slow' }, operationId: 'op-timeout', operationKey: 'demo-plugin.echo', }), ).rejects.toBeInstanceOf(QqbotPluginRuntimeError); await expect( runtime.executeOperation({ input: { text: 'slow' }, operationId: 'op-timeout-2', operationKey: 'demo-plugin.echo', }), ).rejects.toMatchObject({ code: 'PLUGIN_WORKER_TIMEOUT', }); expect(driver.dispose).toHaveBeenCalled(); expect(runtime.status).toBe('failed'); }); it('keeps successful worker calls alive after the timeout window', async () => { const driver: QqbotPluginWorkerDriver = { dispose: jest.fn(async () => undefined), request: jest.fn(async () => ({ ok: true })), }; const runtime = new QqbotPluginWorkerRuntime( new RecordingRequestQueue(driver), { defaultTimeoutMs: 5, installationId: 'install-success', pluginKey: 'demo-plugin', }, ); await runtime.load({ entry: 'src/index.ts', pluginKey: 'demo-plugin', version: '0.1.0', }); await runtime.activate(); await new Promise((resolve) => setTimeout(resolve, 20)); expect(driver.dispose).not.toHaveBeenCalled(); expect(runtime.status).toBe('active'); expect(runtime.listRuntimeEvents()).toEqual([]); }); }); /** * 创建 QQBot 插件平台对象或配置。 */ function createDeferred() { let resolve!: (value: T | PromiseLike) => void; let reject!: (reason?: unknown) => void; const promise = new Promise((innerResolve, innerReject) => { resolve = innerResolve; reject = innerReject; }); return { promise, reject, resolve }; } /** * 执行 QQBot 插件平台流程。 * @param condition - condition 输入;驱动 `Error()` 的 插件平台步骤。 * @param timeoutMs - 插件平台列表;驱动 `Date.now()` 的 插件平台步骤。 */ async function waitUntil( condition: () => boolean, timeoutMs = 100, ): Promise { const deadline = Date.now() + timeoutMs; while (!condition()) { if (Date.now() > deadline) { throw new Error('condition was not met before timeout'); } await new Promise((resolve) => setTimeout(resolve, 1)); } } describe('QQBot plugin SDK contract', () => { it('exposes only host-controlled capabilities to plugin code', () => { const sdk = createQqbotPluginSdk({ assets: { readAsset: jest.fn(), }, config: { getConfig: jest.fn(), setConfig: jest.fn(), }, eventContext: { eventKey: 'demo-plugin.message', }, events: { emitRuntimeEvent: jest.fn(), }, http: { request: jest.fn(), }, operationContext: { operationKey: 'demo-plugin.echo', }, sendQueue: { sendMessage: jest.fn(), }, storage: { get: jest.fn(), set: jest.fn(), }, }); expect(Object.keys(sdk).sort()).toEqual([ 'assets', 'config', 'eventContext', 'events', 'http', 'operationContext', 'sendQueue', 'storage', ]); expect('env' in sdk).toBe(false); expect('fs' in sdk).toBe(false); expect('nest' in sdk).toBe(false); expect('repository' in sdk).toBe(false); }); }); describe('QQBot plugin worker queue config', () => { it('requires an explicit Redis host for plugin worker queues', () => { const configService = createConfigService({}); expect(() => resolveQqbotPluginQueueConnection(configService)).toThrow( 'QQBot 插件队列缺少 Redis 主机配置', ); }); it('uses safe numeric defaults when optional Redis values are blank', () => { const options = createQqbotBullmqWorkerQueueOptions( createConfigService({ QQBOT_PLUGIN_QUEUE_REDIS_DB: '', QQBOT_PLUGIN_QUEUE_REDIS_HOST: 'redis.local', QQBOT_PLUGIN_QUEUE_REDIS_PORT: '', QQBOT_PLUGIN_QUEUE_REDIS_PREFIX: '', QQBOT_PLUGIN_QUEUE_REMOVE_ON_FAIL: '', QQBOT_PLUGIN_QUEUE_WAIT_BUFFER_MS: '', }), 'bangdream', 'install-1', ); expect(options).toMatchObject({ connection: { db: 0, host: 'redis.local', port: 6379, }, installationId: 'install-1', pluginKey: 'bangdream', prefix: 'kt:qqbot:plugin-worker', queueWaitTimeoutMs: 120_000, removeOnFailCount: 100, waitUntilFinishedBufferMs: 5_000, workerInstanceId: expect.any(String), }); expect(options.workerInstanceId).toContain(':'); }); }); /** * 创建 QQBot 插件平台对象或配置。 * @param values - 配置值字典;生成 插件平台对象。 * @returns 创建后的 QQBot 插件平台对象或配置。 */ function createConfigService(values: Record): ConfigService { return { get: (key: string) => values[key] as T, } as ConfigService; }