]> Piment Noir Git Repositories - e-mobility-charging-stations-simulator.git/commitdiff
feat(worker): terminate a deleted station's worker once it hosts zero elements (...
authorJérôme Benoit <jerome.benoit@piment-noir.org>
Sun, 19 Jul 2026 19:17:18 +0000 (21:17 +0200)
committerGitHub <noreply@github.com>
Sun, 19 Jul 2026 19:17:18 +0000 (21:17 +0200)
* feat(worker): terminate a deleted station's worker once it hosts zero elements (#2027)

Add an element-granular removeElement primitive to the worker abstraction and
wire it into the station-delete path so a deleted charging station's hosting
worker thread is terminated once it no longer hosts any station, without
disrupting sibling stations that share the worker (elementsPerWorker > 1).

- WorkerAbstract: new abstract removeElement(elementKey: PropertyKey).
- WorkerSet: internal PropertyKey -> WorkerSetElement map fed by an injected
  generic elementKey selector; decrement numberOfWorkerElements and reuse a
  factored terminateWorker helper (shared with stop()) to terminate only at
  zero elements with no in-flight addition; purge the map on element removal.
- WorkerFixedPool/WorkerDynamicPool: documented no-op (poolifier exposes no safe
  single-element eviction without destroying the worker and its siblings).
- Bootstrap: pass stationInfo => stationInfo.hashId; call removeElement from
  workerEventDeleted.

Set-vs-pool decision: element-granular termination is workerSet-only; pool
worker threads are a bounded, reused resource and are not reclaimed per delete.
Cleanup/timeout contract: termination reuses the existing worker exit handler
for pending-promise rejection and set/map cleanup; any in-flight broadcast
request still resolves via the existing 60s aggregation timeout backstop.
Part A (#2031 cancellable reset) and bulk stop() are untouched.

* fix(worker): harden element-granular termination against concurrency races

Cross-validated review of the removeElement primitive surfaced two reachable
defects and one latent robustness issue, now fixed:

- Phantom-count leak on re-adding the same element key: numberOfWorkerElements
  was incremented unconditionally, so a duplicate/re-added key inflated the
  count and the worker was never terminated. Count is now key-aware (distinct
  keys per worker); factored into trackAddedWorkerElement.
- New addition routed onto a terminating worker: getWorkerSetElement could
  select a worker parked in terminateWorker, losing the message and rejecting
  an unrelated add. Terminating workers are now flagged and excluded from
  selection.
- getWorkerSetElementByWorker matched by threadId, which collapses to -1 after
  terminate(); switched to worker object identity.

Tests: add coverage for the in-flight-sibling gate, drain-to-last-sibling,
add-during-terminate, same-key re-add, and unknown-key no-op; consolidate the
silent worker fixture into echoWorker via a `hold` flag. Mutation-verified.

* fix(worker): make element termination robust and tighten its types

Second cross-validated review round hardened the removeElement primitive:

- terminateWorker no longer awaits a re-attached 'exit' listener (redundant,
  since worker.terminate() fulfils once the worker has exited, and it could
  deadlock if 'exit' fired first). It now catches a rejected terminate()
  (surfaced via the error event, not rethrown so the removal still succeeds)
  and performs pending-promise rejection + set/map cleanup in a finally, closing
  the zombie-on-reject leak, the orphaned-pending edge, and the stop-concurrent-
  exit deadlock in one place.
- Drop the dead, semantically fictional migration branch in
  trackAddedWorkerElement: duplicate live element keys cannot occur (identities
  are deduplicated upstream), so a key is counted once per worker and a same-key
  re-add is a no-op.
- Make WorkerSetElement.terminating a required boolean (initialized false) and
  narrow the element key from PropertyKey to string, matching the sole caller
  (hashId) and avoiding non-serializable keys.

Tests: add reject-path coverage (removeElement resolves and cleans up when
terminate() rejects; an in-flight addition is still rejected on stop when
terminate() rejects). Mutation-verified. All project gates green.

* test(worker): align WorkerSet internals view element key type to string

Coherence with the production narrowing of the element key from PropertyKey
to string; test-only, no runtime change.

* docs(worker): drop the worker-termination note from the README

src/charging-station/Bootstrap.ts
src/worker/WorkerAbstract.ts
src/worker/WorkerDynamicPool.ts
src/worker/WorkerFactory.ts
src/worker/WorkerFixedPool.ts
src/worker/WorkerSet.ts
src/worker/WorkerTypes.ts
tests/worker/WorkerSet.test.ts [new file with mode: 0644]
tests/worker/fixtures/echoWorker.mjs [new file with mode: 0644]

index a68b8a05d74b96c11df754fdb9b909487be3ced9..f8702b5a536b9cbc173f32512cd7a5701b824027 100644 (file)
@@ -559,7 +559,8 @@ export class Bootstrap extends EventEmitter implements IBootstrap {
           }),
         },
         workerStartDelay: workerConfiguration.startDelay,
-      }
+      },
+      (stationInfo: ChargingStationInfo) => stationInfo.hashId
     )
   }
 
@@ -746,6 +747,15 @@ export class Bootstrap extends EventEmitter implements IBootstrap {
   }
 
   private readonly workerEventDeleted = (data: ChargingStationData): void => {
+    this.workerImplementation?.removeElement(data.stationInfo.hashId).catch((error: unknown) => {
+      logger.error(
+        `${this.logPrefix()} ${moduleName}.workerEventDeleted: Error while terminating the emptied worker of charging station ${
+          // eslint-disable-next-line @typescript-eslint/restrict-template-expressions
+          data.stationInfo.chargingStationId
+        } (hashId: ${data.stationInfo.hashId}):`,
+        error
+      )
+    })
     if (this.uiServer.deleteChargingStationData(data.stationInfo.hashId)) {
       this.uiServer.scheduleClientNotification()
     }
index a2946a1f8e6533bed7c5cb2a98ce61db9f1c4c27..a8dcb0f66b9b8706f33b8dac966492c80545f6ad 100644 (file)
@@ -45,6 +45,13 @@ export abstract class WorkerAbstract<D extends WorkerData, R extends WorkerData>
    * @param elementData - The element data to process
    */
   public abstract addElement (elementData: D): Promise<R>
+  /**
+   * Removes a task element from the worker pool/set, terminating its hosting
+   * worker once that worker no longer hosts any element. Sibling elements
+   * sharing the worker are left running.
+   * @param elementKey - Identity key of the element to remove
+   */
+  public abstract removeElement (elementKey: string): Promise<void>
   /**
    * Starts the worker pool/set.
    */
index 66febb575be54e7f6192267a0c1c1c887076159e..15ce24013ce2a8d7940cfd094ea034b83084e4a0 100644 (file)
@@ -54,6 +54,17 @@ export class WorkerDynamicPool<D extends WorkerData, R extends WorkerData> exten
     return response
   }
 
+  /**
+   * @inheritDoc
+   * @remarks
+   * No-op: poolifier exposes no safe primitive to evict a single element from a
+   * pool worker without destroying the worker and its siblings. Pool worker
+   * threads are a bounded, reused resource (≤ `poolMaxSize`).
+   */
+  public removeElement (): Promise<void> {
+    return Promise.resolve()
+  }
+
   /** @inheritDoc */
   public start (): void {
     this.pool.start()
index 4807c88dce374c7ba78afb7081a3eeb687ff3d78..eacf6ab0c2394506cb50416a84640c08215653a2 100644 (file)
@@ -18,7 +18,8 @@ export class WorkerFactory {
   public static getWorkerImplementation<D extends WorkerData, R extends WorkerData>(
     workerScript: string,
     workerProcessType: WorkerProcessType,
-    workerOptions?: WorkerOptions
+    workerOptions?: WorkerOptions,
+    elementKey?: (element: R) => string
   ): WorkerAbstract<D, R> {
     if (!isMainThread) {
       throw new Error('Cannot get a worker implementation outside the main thread')
@@ -30,7 +31,7 @@ export class WorkerFactory {
       case WorkerProcessType.fixedPool:
         return new WorkerFixedPool<D, R>(workerScript, resolvedOptions)
       case WorkerProcessType.workerSet:
-        return new WorkerSet<D, R>(workerScript, resolvedOptions)
+        return new WorkerSet<D, R>(workerScript, resolvedOptions, elementKey)
       default:
         // eslint-disable-next-line @typescript-eslint/restrict-template-expressions
         throw new Error(`Worker implementation type '${workerProcessType}' not found`)
index c05864d1b78fa2d951c2c1f33d59b1e02a5b1304..f988e3e9aaedf15680d409a9cae7723822de8b35 100644 (file)
@@ -53,6 +53,17 @@ export class WorkerFixedPool<D extends WorkerData, R extends WorkerData> extends
     return response
   }
 
+  /**
+   * @inheritDoc
+   * @remarks
+   * No-op: poolifier exposes no safe primitive to evict a single element from a
+   * pool worker without destroying the worker and its siblings. Pool worker
+   * threads are a bounded, reused resource (≤ `poolMaxSize`).
+   */
+  public removeElement (): Promise<void> {
+    return Promise.resolve()
+  }
+
   /** @inheritDoc */
   public start (): void {
     this.pool.start()
index 6ee6ca6cc46e52fcd97898fb708f232335da5e0a..f08b0d02ec67a3458917ca7f5b731a69967a1f36 100644 (file)
@@ -54,6 +54,8 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
     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
@@ -64,8 +66,14 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
    * 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')
@@ -78,6 +86,8 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
     }
     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' })
     }
@@ -112,6 +122,27 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
     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()
@@ -126,15 +157,7 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
   /** @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 {
@@ -178,7 +201,7 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
         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
@@ -237,6 +260,7 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
     })
     const workerSetElement: WorkerSetElement = {
       numberOfWorkerElements: 0,
+      terminating: false,
       worker,
     }
     this.workerSet.add(workerSetElement)
@@ -248,8 +272,9 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
     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
@@ -268,7 +293,7 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
   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
       }
@@ -276,6 +301,15 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
     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) {
@@ -294,6 +328,11 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
     if (workerSetElement == null) {
       return
     }
+    for (const [key, mappedWorkerSetElement] of this.elementMap) {
+      if (mappedWorkerSetElement === workerSetElement) {
+        this.elementMap.delete(key)
+      }
+    }
     this.workerSet.delete(workerSetElement)
   }
 
@@ -302,4 +341,39 @@ export class WorkerSet<D extends WorkerData, R extends WorkerData> extends Worke
       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)
+  }
 }
index e6324f5c5cfa51ca421c231847d35aafe745aa89..d23a82b675122f0bdd5f87e3799a9f7f706fccd0 100644 (file)
@@ -67,5 +67,6 @@ export interface WorkerOptions extends Record<string, unknown> {
 
 export interface WorkerSetElement {
   numberOfWorkerElements: number
+  terminating: boolean
   worker: Worker
 }
diff --git a/tests/worker/WorkerSet.test.ts b/tests/worker/WorkerSet.test.ts
new file mode 100644 (file)
index 0000000..87bcdb7
--- /dev/null
@@ -0,0 +1,240 @@
+/**
+ * @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()
+  })
+})
diff --git a/tests/worker/fixtures/echoWorker.mjs b/tests/worker/fixtures/echoWorker.mjs
new file mode 100644 (file)
index 0000000..166f036
--- /dev/null
@@ -0,0 +1,13 @@
+import { parentPort } from 'node:worker_threads'
+
+// Minimal standalone WorkerSet element worker used by WorkerSet.test.ts: echoes
+// each addWorkerElement request back as an addedWorkerElement response so the
+// set can track its per-worker element count, and stays alive until terminated.
+// A request whose data carries `hold: true` is deliberately left unanswered, to
+// exercise in-flight-addition handling.
+parentPort?.on('message', message => {
+  const { data, event, uuid } = message
+  if (event === 'addWorkerElement' && data?.hold !== true) {
+    parentPort?.postMessage({ data, event: 'addedWorkerElement', uuid })
+  }
+})