return this.workerSet.size
}
+ private readonly elementKey?: (element: R) => string
+ private readonly elementMap: Map<string, WorkerSetElement>
private readonly promiseResponseMap: Map<UUIDv4, ResponseWrapper<R>>
private started: boolean
* Creates a new `WorkerSet`.
* @param workerScript - Path to the worker script file
* @param workerOptions - Worker set configuration options
+ * @param elementKey - Optional selector deriving an element's identity key
+ * from its `addElement` response, enabling element-granular removal
*/
- constructor (workerScript: string, workerOptions: WorkerOptions) {
+ constructor (
+ workerScript: string,
+ workerOptions: WorkerOptions,
+ elementKey?: (element: R) => string
+ ) {
super(workerScript, workerOptions)
if (this.workerOptions.elementsPerWorker == null) {
throw new TypeError('Elements per worker is not defined')
}
this.workerSet = new Set<WorkerSetElement>()
this.promiseResponseMap = new Map<UUIDv4, ResponseWrapper<R>>()
+ this.elementMap = new Map<string, WorkerSetElement>()
+ this.elementKey = elementKey
if (this.workerOptions.poolOptions?.enableEvents === true) {
this.emitter = new EventEmitterAsyncResource({ name: 'workerset' })
}
return response
}
+ /** @inheritDoc */
+ public async removeElement (elementKey: string): Promise<void> {
+ const workerSetElement = this.elementMap.get(elementKey)
+ if (workerSetElement == null) {
+ return
+ }
+ this.elementMap.delete(elementKey)
+ if (workerSetElement.numberOfWorkerElements > 0) {
+ --workerSetElement.numberOfWorkerElements
+ }
+ // Terminate only at zero elements with no in-flight addition, so siblings
+ // are never disrupted. terminateWorker owns the pending-promise rejection
+ // and set/map cleanup.
+ if (
+ workerSetElement.numberOfWorkerElements === 0 &&
+ !this.hasPendingElementForWorker(workerSetElement)
+ ) {
+ await this.terminateWorker(workerSetElement)
+ }
+ }
+
/** @inheritDoc */
public async start (): Promise<void> {
this.addWorkerSetElement()
/** @inheritDoc */
public async stop (): Promise<void> {
for (const workerSetElement of this.workerSet) {
- const worker = workerSetElement.worker
- const waitWorkerExit = new Promise<void>(resolve => {
- worker.once('exit', () => {
- resolve()
- })
- })
- worker.unref()
- await worker.terminate()
- await waitWorkerExit
+ await this.terminateWorker(workerSetElement)
}
for (const [uuid, responseWrapper] of this.promiseResponseMap) {
try {
const { reject, resolve, workerSetElement } = responseWrapper
switch (event) {
case WorkerMessageEvents.addedWorkerElement:
- ++workerSetElement.numberOfWorkerElements
+ this.trackAddedWorkerElement(workerSetElement, data)
this.safeEmit(WorkerSetEvents.elementAdded, this.info)
resolve(data)
break
})
const workerSetElement: WorkerSetElement = {
numberOfWorkerElements: 0,
+ terminating: false,
worker,
}
this.workerSet.add(workerSetElement)
let chosenWorkerSetElement: undefined | WorkerSetElement
for (const workerSetElement of this.workerSet) {
if (
+ !workerSetElement.terminating &&
workerSetElement.numberOfWorkerElements <
- (this.workerOptions.elementsPerWorker ?? DEFAULT_ELEMENTS_PER_WORKER)
+ (this.workerOptions.elementsPerWorker ?? DEFAULT_ELEMENTS_PER_WORKER)
) {
chosenWorkerSetElement = workerSetElement
break
private getWorkerSetElementByWorker (worker: Worker): undefined | WorkerSetElement {
let workerSetElementFound: undefined | WorkerSetElement
for (const workerSetElement of this.workerSet) {
- if (workerSetElement.worker.threadId === worker.threadId) {
+ if (workerSetElement.worker === worker) {
workerSetElementFound = workerSetElement
break
}
return workerSetElementFound
}
+ private hasPendingElementForWorker (workerSetElement: WorkerSetElement): boolean {
+ for (const responseWrapper of this.promiseResponseMap.values()) {
+ if (responseWrapper.workerSetElement === workerSetElement) {
+ return true
+ }
+ }
+ return false
+ }
+
private rejectPendingPromiseForWorker (workerSetElement: WorkerSetElement, reason: unknown): void {
for (const [uuid, responseWrapper] of this.promiseResponseMap) {
if (responseWrapper.workerSetElement === workerSetElement) {
if (workerSetElement == null) {
return
}
+ for (const [key, mappedWorkerSetElement] of this.elementMap) {
+ if (mappedWorkerSetElement === workerSetElement) {
+ this.elementMap.delete(key)
+ }
+ }
this.workerSet.delete(workerSetElement)
}
this.emitter.emit(event, ...args)
}
}
+
+ // terminate() fulfils once the worker has exited, so waiting on the 'exit'
+ // event afterwards is redundant and would deadlock if the event fired before
+ // its listener was attached. A rejected terminate() is surfaced but not
+ // rethrown: the element removal still succeeds. Cleanup is idempotent with the
+ // 'exit' handler.
+ private async terminateWorker (workerSetElement: WorkerSetElement): Promise<void> {
+ workerSetElement.terminating = true
+ workerSetElement.worker.unref()
+ try {
+ await workerSetElement.worker.terminate()
+ } catch (error) {
+ this.safeEmit(WorkerSetEvents.error, error)
+ } finally {
+ this.rejectPendingPromiseForWorker(workerSetElement, new Error('Worker terminated'))
+ this.removeWorkerSetElement(workerSetElement)
+ }
+ }
+
+ // A worker hosts each distinct element key at most once, so count only new
+ // keys; re-adding a key already tracked on the same worker is a no-op.
+ // Duplicate live keys cannot land on different workers (identities are
+ // deduplicated upstream). Without a selector the count is a plain tally.
+ private trackAddedWorkerElement (workerSetElement: WorkerSetElement, element: R): void {
+ if (this.elementKey == null) {
+ ++workerSetElement.numberOfWorkerElements
+ return
+ }
+ const key = this.elementKey(element)
+ if (this.elementMap.get(key) === workerSetElement) {
+ return
+ }
+ ++workerSetElement.numberOfWorkerElements
+ this.elementMap.set(key, workerSetElement)
+ }
}
--- /dev/null
+/**
+ * @file Tests for WorkerSet
+ * @description Element-granular worker termination (issue #2027): a removed
+ * element's hosting worker is terminated once it hosts zero elements, without
+ * disrupting sibling elements or unrelated in-flight additions, and bulk stop()
+ * stays unchanged with no leaked pending-response or element-map entries.
+ */
+import assert from 'node:assert/strict'
+import { afterEach, describe, it } from 'node:test'
+import { fileURLToPath } from 'node:url'
+
+import type { WorkerData } from '../../src/worker/index.js'
+
+import { WorkerSet } from '../../src/worker/WorkerSet.js'
+import { flushMicrotasks, standardCleanup } from '../helpers/TestLifecycleHelpers.js'
+
+interface TestElement extends WorkerData {
+ hold?: boolean
+ id: string
+}
+
+interface WorkerSetElementView {
+ numberOfWorkerElements: number
+ terminating: boolean
+ worker: { terminate: () => Promise<number>; threadId: number }
+}
+
+interface WorkerSetInternals {
+ elementMap: Map<string, WorkerSetElementView>
+ promiseResponseMap: Map<unknown, unknown>
+ workerSet: Set<WorkerSetElementView>
+}
+
+const ECHO_WORKER_SCRIPT = fileURLToPath(new URL('./fixtures/echoWorker.mjs', import.meta.url))
+
+const createWorkerSet = (elementsPerWorker: number): WorkerSet<TestElement, TestElement> =>
+ new WorkerSet<TestElement, TestElement>(
+ ECHO_WORKER_SCRIPT,
+ {
+ elementAddDelay: 0,
+ elementsPerWorker,
+ poolMaxSize: 4,
+ poolMinSize: 1,
+ workerStartDelay: 0,
+ },
+ element => element.id
+ )
+
+const internalsOf = (workerSet: WorkerSet<TestElement, TestElement>): WorkerSetInternals =>
+ workerSet as unknown as WorkerSetInternals
+
+await describe('WorkerSet', async () => {
+ let workerSet: undefined | WorkerSet<TestElement, TestElement>
+
+ afterEach(async () => {
+ await workerSet?.stop()
+ workerSet = undefined
+ standardCleanup()
+ })
+
+ await it('should terminate the hosting worker once it hosts zero elements', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ const info = await workerSet.addElement({ id: 'cs-1' })
+ assert.strictEqual(info.id, 'cs-1')
+ assert.strictEqual(workerSet.size, 1)
+
+ await workerSet.removeElement('cs-1')
+
+ assert.strictEqual(workerSet.size, 0)
+ })
+
+ await it('should not terminate a worker still hosting sibling elements', async () => {
+ workerSet = createWorkerSet(2)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ await workerSet.addElement({ id: 'cs-2' })
+ assert.strictEqual(workerSet.size, 1)
+ const internals = internalsOf(workerSet)
+ const [hostingWorker] = [...internals.workerSet]
+ const threadIdBefore = hostingWorker.worker.threadId
+
+ await workerSet.removeElement('cs-1')
+
+ // The shared worker survives with only the sibling element remaining.
+ assert.strictEqual(workerSet.size, 1)
+ assert.strictEqual([...internals.workerSet][0].worker.threadId, threadIdBefore)
+ assert.strictEqual(hostingWorker.numberOfWorkerElements, 1)
+ assert.strictEqual(internals.elementMap.has('cs-2'), true)
+ assert.strictEqual(internals.elementMap.has('cs-1'), false)
+ })
+
+ await it('should terminate only at the last sibling when draining a shared worker', async () => {
+ workerSet = createWorkerSet(3)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ await workerSet.addElement({ id: 'cs-2' })
+ await workerSet.addElement({ id: 'cs-3' })
+ assert.strictEqual(workerSet.size, 1)
+
+ await workerSet.removeElement('cs-1')
+ assert.strictEqual(workerSet.size, 1)
+ await workerSet.removeElement('cs-2')
+ assert.strictEqual(workerSet.size, 1)
+ await workerSet.removeElement('cs-3')
+ assert.strictEqual(workerSet.size, 0)
+ })
+
+ await it('should not disrupt a sibling element addition still in flight', async () => {
+ workerSet = createWorkerSet(2)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ const pendingSibling = workerSet.addElement({ hold: true, id: 'cs-2' })
+ await flushMicrotasks()
+ assert.strictEqual(internalsOf(workerSet).promiseResponseMap.size, 1)
+
+ await workerSet.removeElement('cs-1')
+
+ // The in-flight sibling keeps the worker alive despite the empty count.
+ assert.strictEqual(workerSet.size, 1)
+ await workerSet.stop()
+ await assert.rejects(pendingSibling)
+ })
+
+ await it('should not route a new element addition onto a terminating worker', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ // Start terminating cs-1's worker, then add an unrelated element in the same
+ // window: it must land on a fresh worker, not the terminating one.
+ const removing = workerSet.removeElement('cs-1')
+ const info = await workerSet.addElement({ id: 'cs-9' })
+ await removing
+
+ assert.strictEqual(info.id, 'cs-9')
+ assert.strictEqual(workerSet.size, 1)
+ assert.strictEqual(internalsOf(workerSet).elementMap.has('cs-9'), true)
+ })
+
+ await it('should keep the count consistent when the same element key is re-added', async () => {
+ workerSet = createWorkerSet(2)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ await workerSet.addElement({ id: 'cs-1' })
+ const [hostingWorker] = [...internalsOf(workerSet).workerSet]
+ assert.strictEqual(hostingWorker.numberOfWorkerElements, 1)
+
+ await workerSet.removeElement('cs-1')
+
+ // A duplicate add must not inflate the count and orphan the worker.
+ assert.strictEqual(workerSet.size, 0)
+ })
+
+ await it('should no-op when removing an unknown element key', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+
+ await workerSet.removeElement('unknown')
+
+ assert.strictEqual(workerSet.size, 1)
+ assert.strictEqual(internalsOf(workerSet).elementMap.has('cs-1'), true)
+ })
+
+ await it('should terminate all workers on bulk stop', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ await workerSet.addElement({ id: 'cs-2' })
+ assert.strictEqual(workerSet.size, 2)
+
+ await workerSet.stop()
+
+ assert.strictEqual(workerSet.size, 0)
+ assert.strictEqual(workerSet.info.started, false)
+ assert.strictEqual(internalsOf(workerSet).promiseResponseMap.size, 0)
+ })
+
+ await it('should not leak pending-response or element-map entries after termination', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ const internals = internalsOf(workerSet)
+ assert.strictEqual(internals.promiseResponseMap.size, 0)
+ assert.strictEqual(internals.elementMap.size, 1)
+
+ await workerSet.removeElement('cs-1')
+
+ assert.strictEqual(internals.elementMap.size, 0)
+ assert.strictEqual(internals.promiseResponseMap.size, 0)
+ assert.strictEqual(workerSet.size, 0)
+ })
+
+ await it('should reject an in-flight element addition when its worker is terminated', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ const pendingAddition = workerSet.addElement({ hold: true, id: 'cs-1' })
+ await flushMicrotasks()
+ assert.strictEqual(internalsOf(workerSet).promiseResponseMap.size, 1)
+
+ await workerSet.stop()
+
+ await assert.rejects(pendingAddition)
+ })
+
+ await it('should resolve and clean up when worker.terminate() rejects', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ await workerSet.addElement({ id: 'cs-1' })
+ const [hostingWorker] = [...internalsOf(workerSet).workerSet]
+ const realTerminate = hostingWorker.worker.terminate.bind(hostingWorker.worker)
+ hostingWorker.worker.terminate = () => Promise.reject(new Error('ERR_WORKER_NOT_RUNNING'))
+
+ // A rejected terminate() must not hang or reject removeElement; the element
+ // is still removed and the worker record cleaned up.
+ await workerSet.removeElement('cs-1')
+
+ assert.strictEqual(workerSet.size, 0)
+ assert.strictEqual(internalsOf(workerSet).elementMap.size, 0)
+ await realTerminate()
+ })
+
+ await it('should reject an in-flight addition when its worker terminate() rejects on stop', async () => {
+ workerSet = createWorkerSet(1)
+ await workerSet.start()
+ const pendingAddition = workerSet.addElement({ hold: true, id: 'cs-1' })
+ await flushMicrotasks()
+ const [hostingWorker] = [...internalsOf(workerSet).workerSet]
+ const realTerminate = hostingWorker.worker.terminate.bind(hostingWorker.worker)
+ hostingWorker.worker.terminate = () => Promise.reject(new Error('ERR_WORKER_NOT_RUNNING'))
+
+ await workerSet.stop()
+
+ // terminateWorker's cleanup rejects the pending addition even though the
+ // worker never emitted 'exit'.
+ await assert.rejects(pendingAddition)
+ assert.strictEqual(internalsOf(workerSet).promiseResponseMap.size, 0)
+ await realTerminate()
+ })
+})