From e42458a21c8b164467aa66e7da8e17a0d631c5bb Mon Sep 17 00:00:00 2001 From: qnbs <155236708+qnbs@users.noreply.github.com> Date: Sat, 1 Aug 2026 02:41:00 +0200 Subject: [PATCH 01/12] fix(worker-bus): enforce runOnPool() timeoutMs with watchdog + force-terminate - runOnPool() now arms a timer keyed to task.timeoutMs when the TASK message is posted. Previously a wedged worker (no crash, no message) left the task's result promise pending forever. - PROGRESS messages re-arm the watchdog so long-running jobs (ProForge stages, LoRA training) that report progress aren't falsely killed. - On expiry: settle with a recoverable TaskResult error (code TIMEOUT), which flows through the existing retry/circuit-breaker/DLQ pipeline unchanged (recordFailure(), retry-then-DLQ-then-reject). - Add WorkerPool.terminateWorker(workerId): force-terminates one worker and respawns a replacement (mirrors the crash-detection restart path) instead of recycling a wedged worker back to idle via release(). - Tests: timeout rejects a hung task, force-terminates the worker, and is not called for release(); watchdog reset via PROGRESS keeps a slow-but-alive task from being killed; terminateWorker unit tests. --- packages/worker-bus/src/workerBus.ts | 55 +++++++++- packages/worker-bus/src/workerPool.ts | 12 +++ packages/worker-bus/tests/workerBus.test.ts | 101 +++++++++++++++++++ packages/worker-bus/tests/workerPool.test.ts | 37 +++++++ 4 files changed, 200 insertions(+), 5 deletions(-) diff --git a/packages/worker-bus/src/workerBus.ts b/packages/worker-bus/src/workerBus.ts index f7e4b964e..9bb7fee8b 100644 --- a/packages/worker-bus/src/workerBus.ts +++ b/packages/worker-bus/src/workerBus.ts @@ -262,6 +262,9 @@ export class WorkerBus { const worker = await pool.acquire(token.signal); const startedAt = performance.now(); const queueTimeMs = Math.round(startedAt - task.createdAt); + // QNBS-v3: [P0 timeout enforcement — runOnPool previously awaited RESULT/abort forever; + // a wedged worker (no crash, no message) left the task's promise pending indefinitely.] + let timedOut = false; try { const port = worker.port; @@ -271,12 +274,49 @@ export class WorkerBus { const result = await new Promise>((resolve, _reject) => { let settled = false; + let timeoutTimer: ReturnType; + + const settle = () => { + settled = true; + clearTimeout(timeoutTimer); + port.removeEventListener('message', handler); + }; + + // QNBS-v3: watchdog, not a hard ceiling — re-armed on every PROGRESS message so + // long-running jobs (ProForge stages, LoRA training) that report progress + // aren't falsely killed while genuinely wedged workers still get caught. + const armTimeout = () => { + clearTimeout(timeoutTimer); + timeoutTimer = setTimeout(onTimeout, task.timeoutMs); + }; + + const onTimeout = () => { + if (settled) return; + settle(); + timedOut = true; + resolve({ + taskId: task.taskId, + success: false, + error: { + code: 'TIMEOUT', + message: `Task exceeded its ${task.timeoutMs}ms deadline with no response from the worker`, + recoverable: true, + retryCount: 0, + }, + latencyMs: Math.round(performance.now() - startedAt), + queueTimeMs, + workerId: worker.workerId, + layer: 'web', + }); + }; + const handler = (event: MessageEvent) => { if (settled) return; const msg = validateWorkerMessage(event.data); if (!msg) return; if (msg.kind === 'PROGRESS') { + armTimeout(); this.progress.emit(task.taskId, { taskId: task.taskId, taskType: task.taskType, @@ -286,8 +326,7 @@ export class WorkerBus { timestamp: Date.now(), }); } else if (msg.kind === 'RESULT') { - settled = true; - port.removeEventListener('message', handler); + settle(); const latencyMs = Math.round(performance.now() - startedAt); resolve({ taskId: task.taskId, @@ -309,11 +348,11 @@ export class WorkerBus { } }; port.addEventListener('message', handler); + armTimeout(); const onAbort = () => { if (settled) return; - settled = true; - port.removeEventListener('message', handler); + settle(); port.postMessage(createCancelMessage(task.taskId, 'Aborted')); resolve({ taskId: task.taskId, @@ -342,7 +381,13 @@ export class WorkerBus { } return result; } finally { - pool.release(worker); + // QNBS-v3: a timed-out worker is presumed wedged — force-terminate + respawn instead of + // releasing it back to the idle pool where a future task would reuse it. + if (timedOut) { + pool.terminateWorker(worker.workerId); + } else { + pool.release(worker); + } } } diff --git a/packages/worker-bus/src/workerPool.ts b/packages/worker-bus/src/workerPool.ts index b1d2531ba..88cf0ada8 100644 --- a/packages/worker-bus/src/workerPool.ts +++ b/packages/worker-bus/src/workerPool.ts @@ -82,6 +82,18 @@ export class WorkerPool { this.entries = []; } + /** + * Force-terminate one worker and respawn a replacement (pool stays at capacity). + * QNBS-v3: a task that missed its timeoutMs deadline may have wedged the worker in an + * unknown state — treat it like a crash instead of recycling it back to idle via release(). + */ + terminateWorker(workerId: string): void { + const entry = this.entries.find((e) => e.instance.workerId === workerId); + if (!entry) return; + this.setCrashed(workerId); + this.restartWorker(entry); + } + getHealth(): { totalWorkers: number; idleWorkers: number; diff --git a/packages/worker-bus/tests/workerBus.test.ts b/packages/worker-bus/tests/workerBus.test.ts index c0c83332d..6c64d04cc 100644 --- a/packages/worker-bus/tests/workerBus.test.ts +++ b/packages/worker-bus/tests/workerBus.test.ts @@ -626,4 +626,105 @@ describe('WorkerBus', () => { handle.cancel('test-abort'); await expect(handle.result).rejects.toThrow('cancelled'); }); + + it('times out a hung worker, records a circuit-breaker failure, and force-terminates it', async () => { + // Worker never sends PROGRESS or RESULT — simulates a wedged worker (not crashed, just silent). + const mockPort = { + addEventListener: vi.fn(), + removeEventListener: vi.fn(), + postMessage: vi.fn(), + start: vi.fn(), + }; + + const mockWorker = { + workerId: 'mock-worker-hung', + worker: {} as Worker, + channel: { port1: mockPort, port2: mockPort } as unknown as MessageChannel, + port: mockPort as unknown as MessagePort, + state: 'idle' as const, + capabilities: ['inference.text'] as const, + labels: {}, + }; + + const pool = (bus as unknown as { pools: Map }).pools.get('fake')!; + vi.spyOn(pool, 'acquire').mockResolvedValue( + mockWorker as unknown as import('../src/workerPool').PooledWorkerInstance, + ); + const terminateWorkerSpy = vi.spyOn(pool, 'terminateWorker').mockImplementation(() => {}); + const releaseSpy = vi.spyOn(pool, 'release'); + + const handle = bus.enqueue( + 'test.task', + { data: 1 }, + { timeoutMs: 25, retryPolicy: { maxRetries: 0 } }, + ); + + await expect(handle.result).rejects.toThrow(/exceeded its 25ms deadline/i); + expect(terminateWorkerSpy).toHaveBeenCalledWith('mock-worker-hung'); + expect(releaseSpy).not.toHaveBeenCalled(); + expect(bus.getTelemetry().failedTasks).toBe(1); + expect(bus.getTelemetry().deadLetterCount).toBe(1); + expect(bus.getTelemetry().circuitBreakerStates['test.task']).toBeDefined(); + }); + + it('resets the timeout watchdog on PROGRESS so a slow-but-alive task is not killed early', async () => { + const mockPort = { + addEventListener: vi.fn((type: string, handler: EventListener) => { + if (type === 'message') { + // Two PROGRESS pings inside the timeout window, then RESULT after the original + // deadline would have expired — only survives if PROGRESS re-arms the watchdog. + setTimeout(() => { + handler( + new MessageEvent('message', { + data: { kind: 'PROGRESS', taskId: 'mock-task-id', stage: 'step1', progress: 0.3 }, + }), + ); + }, 15); + setTimeout(() => { + handler( + new MessageEvent('message', { + data: { kind: 'PROGRESS', taskId: 'mock-task-id', stage: 'step2', progress: 0.6 }, + }), + ); + }, 30); + setTimeout(() => { + handler( + new MessageEvent('message', { + data: { + kind: 'RESULT', + taskId: 'mock-task-id', + success: true, + result: 'slow-but-done', + latencyMs: 45, + }, + }), + ); + }, 45); + } + }), + removeEventListener: vi.fn(), + postMessage: vi.fn(), + start: vi.fn(), + }; + + const mockWorker = { + workerId: 'mock-worker-slow', + worker: {} as Worker, + channel: { port1: mockPort, port2: mockPort } as unknown as MessageChannel, + port: mockPort as unknown as MessagePort, + state: 'idle' as const, + capabilities: ['inference.text'] as const, + labels: {}, + }; + + const pool = (bus as unknown as { pools: Map }).pools.get('fake')!; + vi.spyOn(pool, 'acquire').mockResolvedValue( + mockWorker as unknown as import('../src/workerPool').PooledWorkerInstance, + ); + + // timeoutMs (20ms) is shorter than the total run (45ms), but each PROGRESS arrives + // well within 20ms of the previous reset, so the watchdog never fires. + const handle = bus.enqueue('test.task', { data: 1 }, { timeoutMs: 20 }); + await expect(handle.result).resolves.toBe('slow-but-done'); + }); }); diff --git a/packages/worker-bus/tests/workerPool.test.ts b/packages/worker-bus/tests/workerPool.test.ts index af7f99850..89488174c 100644 --- a/packages/worker-bus/tests/workerPool.test.ts +++ b/packages/worker-bus/tests/workerPool.test.ts @@ -204,4 +204,41 @@ describe('WorkerPool', () => { pool.release(worker); expect(pool.getHealth().totalWorkers).toBe(0); }); + + it('terminateWorker force-terminates a specific worker and respawns a replacement', async () => { + const pool = new WorkerPool('test-pool', ['inference.text'], { + maxWorkers: 1, + minWorkers: 1, + idleTimeoutMs: 120_000, + workerScript: '/mock.worker.js', + capabilities: ['inference.text'], + labels: {}, + }); + const worker = await pool.acquire(); + const mockWorker = worker.worker as unknown as MockWorker; + + pool.terminateWorker(worker.workerId); + + expect(mockWorker.terminated).toBe(true); + // QNBS-v3: like restartWorker() on crash — the pool respawns to stay at capacity, + // but the new worker must be a distinct instance, not the wedged one. + expect(pool.getHealth().totalWorkers).toBe(1); + expect(pool.getHealth().crashedWorkers).toBe(0); + await pool.terminateAll(); + }); + + it('terminateWorker is a no-op for an unknown workerId', async () => { + const pool = new WorkerPool('test-pool', ['inference.text'], { + maxWorkers: 1, + minWorkers: 1, + idleTimeoutMs: 120_000, + workerScript: '/mock.worker.js', + capabilities: ['inference.text'], + labels: {}, + }); + await pool.acquire(); + expect(() => pool.terminateWorker('nonexistent-worker-id')).not.toThrow(); + expect(pool.getHealth().totalWorkers).toBe(1); + await pool.terminateAll(); + }); }); From 1a9b13ddb125231167d611861d00def93899c252 Mon Sep 17 00:00:00 2001 From: qnbs <155236708+qnbs@users.noreply.github.com> Date: Sat, 1 Aug 2026 02:52:10 +0200 Subject: [PATCH 02/12] feat(worker-bus): size inference pool maxWorkers to device memory tier MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - workerBusManager.inferencePoolOptions() now derives maxWorkers from localAiDeviceProfiler.detectMemoryTier() instead of a hardcoded 2: high memory tier -> 3, medium -> 2 (unchanged default), low -> 1. Still capped at MAX_WORKERS_INFERENCE. - Export detectMemoryTier from localAiDeviceProfiler.ts (was module-private) — it's a synchronous, side-effect-free read, so callers that only need pool sizing don't pay for the full async generateDeviceProfile() (WebGPU/WebNN/battery probes). - Falls back to maxWorkers 2 (previous default) if memory-tier detection throws, so pool init can never be blocked by this. - Tests: high/medium/low tier sizing + throw-fallback, exercised via the same ensureInferencePool() re-registration path as the existing memory-safety-cap test (initWorkerBus() itself only calls the mocked registry.register(), not bus.registerPool()). --- services/ai/localAiDeviceProfiler.ts | 5 ++- services/workerBusManager.ts | 25 ++++++++++- tests/unit/workerBusManager.test.ts | 66 ++++++++++++++++++++++++++++ 3 files changed, 93 insertions(+), 3 deletions(-) diff --git a/services/ai/localAiDeviceProfiler.ts b/services/ai/localAiDeviceProfiler.ts index acd397150..35ef2aa63 100644 --- a/services/ai/localAiDeviceProfiler.ts +++ b/services/ai/localAiDeviceProfiler.ts @@ -172,7 +172,10 @@ function detectDirectML( // Memory Tier // --------------------------------------------------------------------------- -function detectMemoryTier(): DeviceCapabilityProfile['memoryTier'] { +// QNBS-v3: exported (previously module-private) so callers that only need a fast, synchronous +// memory-tier read — e.g. workerBusManager sizing the inference pool at init — don't +// have to pay for the full async generateDeviceProfile() (WebGPU/WebNN/battery probes). +export function detectMemoryTier(): DeviceCapabilityProfile['memoryTier'] { const deviceMemory = typeof navigator !== 'undefined' && 'deviceMemory' in navigator ? (navigator as Navigator & { deviceMemory?: number }).deviceMemory diff --git a/services/workerBusManager.ts b/services/workerBusManager.ts index caf9a97a9..ad58ca875 100644 --- a/services/workerBusManager.ts +++ b/services/workerBusManager.ts @@ -72,9 +72,12 @@ async function inferencePoolOptions() { const { MAX_WORKERS_INFERENCE, MIN_WORKERS, WORKER_IDLE_TIMEOUT_MS } = await import( '@domain/worker-bus' ); + // QNBS-v3: [P1 — scale inference worker count to the device's memory tier instead of a fixed + // cap. Each replica loads its own transformers.js pipeline (no cross-replica cache + // sharing), so more replicas only help on devices with RAM headroom to spare.] + const maxWorkers = Math.min(await resolveInferenceMaxWorkers(), MAX_WORKERS_INFERENCE); return { - // QNBS-v3: [Capped below MAX_WORKERS_INFERENCE — each replica loads its own transformers.js pipeline (no cross-replica cache sharing), so 4 concurrent workers could mean 4x the model memory footprint under a burst.] - maxWorkers: Math.min(2, MAX_WORKERS_INFERENCE), + maxWorkers, minWorkers: MIN_WORKERS, idleTimeoutMs: WORKER_IDLE_TIMEOUT_MS, workerScript: new URL('../workers/v2/inference.worker.ts', import.meta.url).href, @@ -83,6 +86,24 @@ async function inferencePoolOptions() { }; } +/** + * QNBS-v3: memory-tier-driven worker count — high:3, medium:2 (previous hardcoded default), + * low:1. Falls back to the previous default (2) if device profiling throws (e.g. non-browser + * test environment) so this can never block or fail pool initialization. + */ +async function resolveInferenceMaxWorkers(): Promise { + try { + const { detectMemoryTier } = await import('./ai/localAiDeviceProfiler'); + const tier = detectMemoryTier(); + if (tier === 'high') return 3; + if (tier === 'low') return 1; + return 2; + } catch (err) { + log.warn('Failed to detect memory tier for inference pool sizing; using default', err); + return 2; + } +} + /** Re-register the 'inference' pool if it was removed via terminatePool() — a no-op if already present. */ async function reRegisterInferencePool(bus: WorkerBus): Promise { if (bus.hasPool('inference')) return; diff --git a/tests/unit/workerBusManager.test.ts b/tests/unit/workerBusManager.test.ts index 2a17c768f..3f15aed2b 100644 --- a/tests/unit/workerBusManager.test.ts +++ b/tests/unit/workerBusManager.test.ts @@ -295,4 +295,70 @@ describe('workerBusManager', () => { ); }); }); + + describe('inference pool sizing by memory tier', () => { + // QNBS-v3: [P1 — inference pool maxWorkers is now derived from + // localAiDeviceProfiler.detectMemoryTier() instead of a hardcoded 2. Verified via + // the same re-registration path as the "memory-safety cap" test above, since + // initWorkerBus() only calls registry.register() (mocked, no-op) — bus.registerPool() + // with the real computed options is only exercised by ensureInferencePool().] + afterEach(() => { + vi.doUnmock('../../services/ai/localAiDeviceProfiler'); + }); + + async function expectMaxWorkersForTier( + detectMemoryTier: () => 'high' | 'medium' | 'low', + expectedMaxWorkers: number, + ): Promise { + vi.doMock('../../services/ai/localAiDeviceProfiler', () => ({ detectMemoryTier })); + const { initWorkerBus, ensureInferencePool } = await import( + '../../services/workerBusManager' + ); + await initWorkerBus(); + mockRegisterPool.mockClear(); + mockHasPool.mockReturnValue(false); + + await ensureInferencePool(); + + expect(mockRegisterPool).toHaveBeenCalledWith( + 'inference', + expect.arrayContaining(['inference.text', 'inference.embed']), + expect.objectContaining({ maxWorkers: expectedMaxWorkers }), + ); + } + + it('sizes maxWorkers to 3 on a high memory tier', async () => { + await expectMaxWorkersForTier(() => 'high', 3); + }); + + it('sizes maxWorkers to 2 on a medium memory tier (unchanged default)', async () => { + await expectMaxWorkersForTier(() => 'medium', 2); + }); + + it('sizes maxWorkers to 1 on a low memory tier', async () => { + await expectMaxWorkersForTier(() => 'low', 1); + }); + + it('falls back to maxWorkers 2 when memory-tier detection throws', async () => { + vi.doMock('../../services/ai/localAiDeviceProfiler', () => ({ + detectMemoryTier: () => { + throw new Error('profiler boom'); + }, + })); + const { initWorkerBus, ensureInferencePool } = await import( + '../../services/workerBusManager' + ); + await initWorkerBus(); + mockRegisterPool.mockClear(); + mockHasPool.mockReturnValue(false); + + await ensureInferencePool(); + + expect(mockRegisterPool).toHaveBeenCalledWith( + 'inference', + expect.arrayContaining(['inference.text', 'inference.embed']), + expect.objectContaining({ maxWorkers: 2 }), + ); + }); + }); }); From 29627ec5960fef3afbfded045e55adf085221bdd Mon Sep 17 00:00:00 2001 From: qnbs <155236708+qnbs@users.noreply.github.com> Date: Sat, 1 Aug 2026 03:10:03 +0200 Subject: [PATCH 03/12] fix(worker-bus): forward-declare handler/onTimeout to satisfy DeepSource no-use-before-define --- packages/worker-bus/src/workerBus.ts | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/packages/worker-bus/src/workerBus.ts b/packages/worker-bus/src/workerBus.ts index 9bb7fee8b..20f367b02 100644 --- a/packages/worker-bus/src/workerBus.ts +++ b/packages/worker-bus/src/workerBus.ts @@ -275,6 +275,12 @@ export class WorkerBus { const result = await new Promise>((resolve, _reject) => { let settled = false; let timeoutTimer: ReturnType; + // QNBS-v3: `handler`/`onTimeout` are forward-declared via `let` (hoisted, no TDZ read) + // because `settle`/`armTimeout` close over them before their real definitions + // run — the three closures are mutually recursive, so a strict top-to-bottom + // `const` chain isn't possible without one forward reference. + let handler: (event: MessageEvent) => void; + let onTimeout: () => void; const settle = () => { settled = true; @@ -290,7 +296,7 @@ export class WorkerBus { timeoutTimer = setTimeout(onTimeout, task.timeoutMs); }; - const onTimeout = () => { + onTimeout = () => { if (settled) return; settle(); timedOut = true; @@ -310,7 +316,7 @@ export class WorkerBus { }); }; - const handler = (event: MessageEvent) => { + handler = (event: MessageEvent) => { if (settled) return; const msg = validateWorkerMessage(event.data); if (!msg) return; From 58a52c105b0d9e311e8edf6891927769ea02ac17 Mon Sep 17 00:00:00 2001 From: qnbs <155236708+qnbs@users.noreply.github.com> Date: Sat, 1 Aug 2026 03:35:28 +0200 Subject: [PATCH 04/12] fix(workerBus): filter stale worker messages by taskId, start MessagePort, address review feedback MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - workerBus.ts: filter PROGRESS/RESULT messages by msg.taskId so a released-and-reused MessagePort can't let a cancelled task's late message reset the watchdog or resolve a new task (CodeRabbit + Copilot, real bug) - workerPool.ts: call channel.port1.start() in spawnWorker — an addEventListener-based MessagePort never dispatches without it (Copilot, real bug); add log.warn parity in terminateWorker (CodeRabbit nitpick) - reword TIMEOUT error message for accuracy (Copilot) - reformat QNBS-v3 comments in localAiDeviceProfiler.ts/workerBusManager.ts to the required single-line format (CodeRabbit nitpick) - workerBus.test.ts: fix two pre-existing tests broken by the new taskId filter (hardcoded mock taskId no longer matched the real generated one), convert the two new watchdog tests to fake timers to remove CI wall-clock flakiness risk, add a regression test for the taskId/MessagePort-reuse fix, add missing QNBS-v3 markers - workerPool.test.ts, workerBusManager.test.ts: QNBS-v3 marker / single-line reformat - docs/DEEPSOURCE-REVIEW-LOOP.md: log JS-0067 as a systemic false-positive for this ES-module codebase (informational-only check, not required for merge) Co-Authored-By: GitHub Copilot (Claude Sonnet 5) --- docs/DEEPSOURCE-REVIEW-LOOP.md | 15 ++ packages/worker-bus/src/workerBus.ts | 5 +- packages/worker-bus/src/workerPool.ts | 3 + packages/worker-bus/tests/workerBus.test.ts | 244 +++++++++++++------ packages/worker-bus/tests/workerPool.test.ts | 1 + services/ai/localAiDeviceProfiler.ts | 4 +- services/workerBusManager.ts | 10 +- tests/unit/workerBusManager.test.ts | 6 +- 8 files changed, 195 insertions(+), 93 deletions(-) diff --git a/docs/DEEPSOURCE-REVIEW-LOOP.md b/docs/DEEPSOURCE-REVIEW-LOOP.md index 573dfcb3e..9a3d12aaf 100644 --- a/docs/DEEPSOURCE-REVIEW-LOOP.md +++ b/docs/DEEPSOURCE-REVIEW-LOOP.md @@ -270,6 +270,21 @@ GitHub App resumes auto-reviewing, run **both** loops: CodeAnt for narrative/AI new service/hook needs its branches tested or patch coverage dips below the ~72% target and the `codecov/patch` check fails. #236's hook+panel needed a dedicated hook test (offline/error/disabled/ stale-apply/dictionary/clear branches) to clear it — same posture as #232's binderDepth test. +- **2026-08-01** — On PR #305, `DeepSource: JavaScript` surfaced **JS-0067** ("unexpected function + declaration in the global scope") on top-level `function`/`async function` declarations in + `services/ai/localAiDeviceProfiler.ts` and `services/workerBusManager.ts` — a pattern used + throughout this repo's ES modules (every service file declares top-level functions; this is + idiomatic, Biome-approved module code, not `