diff --git a/package-lock.json b/package-lock.json index a7fede70..2ef25135 100644 --- a/package-lock.json +++ b/package-lock.json @@ -9,6 +9,7 @@ "@sinclair/typebox": "^0.34.47", "@stoprocent/noble": "^2.3.17", "@timesplinter/pimple": "^2.1.1", + "@timesplinter/sequential-task-queue": "^1.3.1", "ajv": "^8.17.1", "ajv-formats": "^3.0.1", "buttplug": "^3.2.2", @@ -24,7 +25,6 @@ "read-last-lines": "^1.8.0", "reflect-metadata": "^0.1.13", "say": "^0.16.0", - "sequential-task-queue": "^1.2.1", "serialport": "^13.0.0", "socket.io": "^4.8.3", "speaker": "https://github.com/SlvCtrlPlus/node-speaker/releases/download/v0.1.0/speaker-v0.1.0.tgz", @@ -3171,6 +3171,12 @@ "integrity": "sha512-Sy0nqk5480ZukdvEfx9mCmKMJ+jCnjiXePO7Flf+4RK5W5bs1qyGdExQJZ6tg6SoCKQekRvpaW6GNGl/Zf3WeA==", "license": "LGPL-3.0-or-later" }, + "node_modules/@timesplinter/sequential-task-queue": { + "version": "1.3.1", + "resolved": "https://registry.npmjs.org/@timesplinter/sequential-task-queue/-/sequential-task-queue-1.3.1.tgz", + "integrity": "sha512-eyMiLZZCx/1xJbQ3oQAyG9b/TBXKgBjpVfb1G6CgjIlsgpYTtZv/euSdJDEM4giLDKOC85JJnqrt66tTiY6aRQ==", + "license": "MIT" + }, "node_modules/@tybys/wasm-util": { "version": "0.10.3", "resolved": "https://registry.npmjs.org/@tybys/wasm-util/-/wasm-util-0.10.3.tgz", @@ -8243,9 +8249,6 @@ "url": "https://opencollective.com/express" } }, - "node_modules/sequential-task-queue": { - "version": "1.2.1" - }, "node_modules/serialport": { "version": "13.0.0", "license": "MIT", diff --git a/package.json b/package.json index a63890f6..1bca2260 100644 --- a/package.json +++ b/package.json @@ -22,7 +22,7 @@ "read-last-lines": "^1.8.0", "reflect-metadata": "^0.1.13", "say": "^0.16.0", - "sequential-task-queue": "^1.2.1", + "@timesplinter/sequential-task-queue": "^1.3.1", "serialport": "^13.0.0", "socket.io": "^4.8.3", "speaker": "https://github.com/SlvCtrlPlus/node-speaker/releases/download/v0.1.0/speaker-v0.1.0.tgz", diff --git a/src/device/detectedDeviceOfferQueue.ts b/src/device/detectedDeviceOfferQueue.ts new file mode 100644 index 00000000..681ffde2 --- /dev/null +++ b/src/device/detectedDeviceOfferQueue.ts @@ -0,0 +1,152 @@ +import { CancellationToken, sequentialTaskQueueEvents, SequentialTaskQueue } from '@timesplinter/sequential-task-queue'; +import { AnyDevice } from './device.js'; +import { DeviceDetectionInfo } from './deviceManager.js'; +import DeviceOfferRejectedError from './deviceOfferRejectedError.js'; +import Logger from '../logging/Logger.js'; +import { logError } from '../util/error.js'; + +export type OfferResult = + | { successful: true, device: D } + | { successful: false, reason: unknown }; + +type DeviceOffer = (cancellationToken: CancellationToken) => Promise; + +export default class DetectedDeviceOfferQueue +{ + private readonly queues: Map = new Map(); + + private readonly logger: Logger; + + public constructor(logger: Logger) { + this.logger = logger; + } + + private getOrCreateQueue(detectionId: string): SequentialTaskQueue + { + let queue = this.queues.get(detectionId); + + if (queue !== undefined) { + return queue; + } + + queue = new SequentialTaskQueue(); + + queue.on(sequentialTaskQueueEvents.drained, () => { + // A revoked (closed) queue must survive its own drain - it's kept around deliberately + // as a tombstone so a late offer can still see it and reject itself. + if (!queue.isClosed) { + this.queues.delete(detectionId); + } + }); + + this.queues.set(detectionId, queue); + + return queue; + } + + public offer(deviceDetectionInfo: DeviceDetectionInfo, deviceOffer: DeviceOffer): Promise> + { + const detectionId = deviceDetectionInfo.detectionId; + const queue = this.getOrCreateQueue(detectionId); + + if (queue.isClosed) { + return Promise.resolve({ + successful: false, + reason: new DeviceOfferRejectedError('Device is not available anymore for offering'), + }); + } + + const task = queue.push((cancellationToken: CancellationToken) => this.runOffer(deviceOffer, cancellationToken)); + + return Promise.resolve(task.then( + (result: OfferResult): OfferResult => { + if (result.successful) { + // Reject every other still-queued offer for this detection id without them + // ever running, since this device has already been claimed. + this.close(detectionId, new DeviceOfferRejectedError('Device has been claimed by another provider')); + } + + return result; + }, + // Reached either if the offer was cancelled while still queued (never even starting - + // our callback above never ran) or if deviceOffer() itself rejected/threw uncaught - + // translate both the same way. + (reason: unknown): OfferResult => ({ + successful: false, + reason: reason, + }) + )); + } + + private async runOffer( + deviceOffer: DeviceOffer, + cancellationToken: CancellationToken + ): Promise> { + const device = await deviceOffer(cancellationToken); + + if (device instanceof DeviceOfferRejectedError) { + return { successful: false, reason: device }; + } + + if (true !== cancellationToken.cancelled) { + return { successful: true, device }; + } + + // In case this offer lost the race against another offer: close the device and reject the offer with a meaningful reason. + try { + await device.close(); + } catch (e: unknown) { + logError(this.logger, `Failed to close device '${device.getDeviceId}' offered after its queue was cleared`, e); + } + + return { + successful: false, + reason: cancellationToken.reason, + }; + } + + /** + * True if a queue currently exists for this detection id - either genuinely active/in-flight, + * or a closed tombstone left behind by revoke(). Does not distinguish between the two; + * callers that need "is a fresh announce still blocked by a past revoke" must call + * dropIfRevoked() first. + */ + public has(detectionId: string): boolean + { + return this.queues.has(detectionId); + } + + public dropIfRevoked(detectionId: string): void + { + const queue = this.queues.get(detectionId); + + if (queue !== undefined && queue.isClosed) { + this.queues.delete(detectionId); + } + } + + private close(detectionId: string, reason: DeviceOfferRejectedError): void + { + const queue = this.queues.get(detectionId); + + if (undefined !== queue) { + void queue.close(true, reason); + } + + this.queues.delete(detectionId); + } + + public revoke(detectionId: string, reason: DeviceOfferRejectedError): void + { + const queue = this.getOrCreateQueue(detectionId); + + void queue.close(true, reason); + } + + public closeAll(reason: DeviceOfferRejectedError): void + { + for (const detectionId of this.queues.keys()) { + this.close(detectionId, reason); + } + } +} diff --git a/src/device/deviceManager.ts b/src/device/deviceManager.ts index d88047a0..c625cbaa 100644 --- a/src/device/deviceManager.ts +++ b/src/device/deviceManager.ts @@ -1,12 +1,14 @@ import { AnyDevice, DeviceEvent, DeviceNotification } from './device.js'; import EventEmitter from 'events'; -import { SequentialTaskQueue } from 'sequential-task-queue'; +import { SequentialTaskQueue } from '@timesplinter/sequential-task-queue'; import DeviceState from './deviceState.js'; import { setIntervalAsync } from '../util/async.js'; import Logger from '../logging/Logger.js'; import { logError } from '../util/error.js'; import { DeviceId } from './deviceId.js'; import SettingsManager from '../settings/settingsManager.js'; +import DeviceOfferRejectedError from './deviceOfferRejectedError.js'; +import DetectedDeviceOfferQueue, { OfferResult } from './detectedDeviceOfferQueue.js'; export type DeviceDetectionInfo = { type: string; @@ -21,15 +23,17 @@ export enum DeviceManagerEvent { deviceNotification = 'deviceNotification', } -type AcquireResult = - | { successful: true } - | { successful: false, reason: string }; +type DisabledDetectedDevice = { + deviceDetectionInfo: DeviceDetectionInfo; + canonicalId: DeviceId; + deviceReleased: Promise; +}; type DeviceManagerEventMap = { [DeviceManagerEvent.deviceConnected]: [device: AnyDevice]; [DeviceManagerEvent.deviceDisconnected]: [device: AnyDevice]; [DeviceManagerEvent.deviceRefreshed]: [device: AnyDevice]; - [DeviceManagerEvent.deviceDetected]: [deviceInfo: DeviceDetectionInfo]; + [DeviceManagerEvent.deviceDetected]: [deviceDetectionInfo: DeviceDetectionInfo]; [DeviceManagerEvent.deviceNotification]: [device: AnyDevice, notification: DeviceNotification]; } @@ -39,19 +43,13 @@ export default class DeviceManager private readonly logger: Logger; - private readonly detectedDeviceAcquireQueue: Map void }[]> = new Map(); + private readonly offerQueue: DetectedDeviceOfferQueue; private readonly connectedDevices: Map; private readonly settingsManager: SettingsManager; - /** - * Devices whose retry is pending because their known device is disabled; re-announced once - * it gets (re-)enabled, see `onSettingsChanged()`. `canonicalId` is the device's final id - * whose enablement gates the retry (protocols may only learn it during a handshake, so it - * can differ from the map key, the preliminary `deviceInfo.detectionId`). - */ - private readonly pendingRetries: Map }> = new Map(); + private readonly detectedDisabledDevices: Map = new Map(); // Serializes onSettingsChanged() runs so rapid settings changes don't interleave private readonly settingsChangeQueue: SequentialTaskQueue = new SequentialTaskQueue(); @@ -66,130 +64,102 @@ export default class DeviceManager this.logger = logger.child({ name: DeviceManager.name }); this.connectedDevices = connectedDevices; this.settingsManager = settingsManager; + this.offerQueue = new DetectedDeviceOfferQueue(this.logger); } public isDeviceEnabled(deviceId: DeviceId): boolean { return this.settingsManager.getSettings()?.getKnownDeviceById(deviceId)?.enabled ?? true; } - public announceDetectedDevice(deviceInfo: DeviceDetectionInfo): void + public announceDetectedDevice(deviceDetectionInfo: DeviceDetectionInfo): void { - if (this.detectedDeviceAcquireQueue.has(deviceInfo.detectionId)) { - return; - } + this.offerQueue.dropIfRevoked(deviceDetectionInfo.detectionId); - if (this.connectedDevices.has(deviceInfo.detectionId)) { - this.logger.debug(`Device with id '${deviceInfo.detectionId}' is already connected, not announcing it as detected`); + if (this.offerQueue.has(deviceDetectionInfo.detectionId)) { return; } - if (!this.isDeviceEnabled(deviceInfo.detectionId)) { - this.logger.debug(`Device with id '${deviceInfo.detectionId}' is disabled, not announcing it as detected`); - // No connection happened yet, so the detection id doubles as the canonical id here - this.registerPendingRetry(deviceInfo, deviceInfo.detectionId); + if (this.connectedDevices.has(deviceDetectionInfo.detectionId)) { + this.logger.debug(`Device with id '${deviceDetectionInfo.detectionId}' is already connected, not announcing it as detected`); return; } - this.logger.info(`Detected new device with id ${deviceInfo.detectionId}`); - - this.detectedDeviceAcquireQueue.set(deviceInfo.detectionId, []); + this.logger.info(`Detected new device with id ${deviceDetectionInfo.detectionId}`); - const hadListeners = this.eventEmitter.emit(DeviceManagerEvent.deviceDetected, deviceInfo); + const hadListeners = this.eventEmitter.emit(DeviceManagerEvent.deviceDetected, deviceDetectionInfo); if (!hadListeners) { - // no subscribed providers, remove empty list from acquire queue for this device - this.logger.info(`No provider available for detected device with id '${deviceInfo.detectionId}'`); - this.detectedDeviceAcquireQueue.delete(deviceInfo.detectionId); + this.logger.info(`No active/started providers. Handling of device detection with id '${deviceDetectionInfo.detectionId}' not possible`); + } else if (!this.offerQueue.has(deviceDetectionInfo.detectionId)) { + this.logger.info(`No provider could handle detected device with id '${deviceDetectionInfo.detectionId}'`); } } - public revokeDetectedDevice(deviceInfo: DeviceDetectionInfo): void + public revokeDetectedDevice(deviceDetectionInfo: DeviceDetectionInfo): void { // A device that physically disappeared should no longer be retried on re-enable - this.pendingRetries.delete(deviceInfo.detectionId); - this.clearDetectedDeviceAcquireQueue(deviceInfo.detectionId, `Device with id '${deviceInfo.detectionId}' has disappeared`); + this.detectedDisabledDevices.delete(deviceDetectionInfo.detectionId); + this.offerQueue.revoke( + deviceDetectionInfo.detectionId, + new DeviceOfferRejectedError('Device has disappeared') + ); } - public async acquireDetectedDevice(deviceId: DeviceId): Promise + public async offerDevice(deviceDetectionInfo: DeviceDetectionInfo, deviceOffer: () => Promise): Promise> { - return new Promise((resolve) => { - const deviceQueue = this.detectedDeviceAcquireQueue.get(deviceId); + if (this.connectedDevices.has(deviceDetectionInfo.detectionId)) { + return { successful: false, reason: new DeviceOfferRejectedError('Device is already connected') }; + } - if (undefined === deviceQueue) { - resolve({ successful: false, reason: `Device with id '${deviceId}' is not available for claiming` }); - return; - } + const result = await this.offerQueue.offer(deviceDetectionInfo, async (cancellationToken) => { + const device = await deviceOffer(); - // Always add to queue first - deviceQueue.push({ resolve }); + if (true === cancellationToken.cancelled) { + try { + await device.close(); + } catch (e: unknown) { + logError(this.logger, `Failed to close device '${device.getDeviceId}' after its offer was cancelled`, e); + } - // If we're first in line, resolve immediately - if (deviceQueue.length === 1) { - resolve({ successful: true }); + return new DeviceOfferRejectedError('Device offer was cancelled'); } - }); - } - public releaseDetectedDevice(deviceId: DeviceId): void - { - const deviceQueue = this.detectedDeviceAcquireQueue.get(deviceId); + if (!this.isDeviceEnabled(device.getDeviceId)) { + this.logger.info(`Not adding device '${device.getDeviceId}' since it is disabled`); - if (undefined === deviceQueue) { - return; - } + const deviceReleased = device.close() + .catch((e: unknown) => logError(this.logger, `Failed to close disabled device '${device.getDeviceId}'`, e)); - // Release current claimant and hand off the claim to the next waiter - deviceQueue.shift(); + // Keyed by detection id so revokeDetectedDevice() (which only has that id) can drop it + this.detectedDisabledDevices.set(deviceDetectionInfo.detectionId, { deviceDetectionInfo, canonicalId: device.getDeviceId, deviceReleased }); - if (deviceQueue.length === 0) { - this.detectedDeviceAcquireQueue.delete(deviceId); - return; - } - - deviceQueue[0]?.resolve({ successful: true }); - } - - /** - * Registers a fully connected device, unless its known device (identified by the final - * `getDeviceId`) is disabled - then the device is closed and registered for retry instead. - * Returns whether the device was added. - */ - public addDevice(deviceInfo: DeviceDetectionInfo, device: AnyDevice): boolean - { - if (!this.isDeviceEnabled(device.getDeviceId)) { - this.logger.info(`Not adding device '${device.getDeviceId}' since it is disabled`); - - const closingDevice = device.close() - .catch((e: unknown) => logError(this.logger, `Failed to close disabled device '${device.getDeviceId}'`, e)); + return new DeviceOfferRejectedError(`Device with id '${device.getDeviceId}' is currently disabled`); + } - this.registerPendingRetry(deviceInfo, device.getDeviceId, closingDevice); - this.releaseDetectedDevice(deviceInfo.detectionId); + return device; + }); - return false; + if (result.successful) { + this.registerDevice(result.device); } - this.connectedDevices.set(device.getDeviceId, device); + return result; + } - device.on(DeviceEvent.deviceRefreshed, (d) => this.refreshDevice(d)); - device.on(DeviceEvent.deviceDisconnected, (d) => this.removeDevice(d)); + private registerDevice(device: AnyDevice): void + { + device.on(DeviceEvent.deviceRefreshed, (d) => this.eventEmitter.emit(DeviceManagerEvent.deviceRefreshed, d)); + device.on(DeviceEvent.deviceDisconnected, (d) => { + this.connectedDevices.delete(d.getDeviceId); + this.eventEmitter.emit(DeviceManagerEvent.deviceDisconnected, d); + }); device.on(DeviceEvent.deviceNotification, (d, notification) => this.eventEmitter.emit(DeviceManagerEvent.deviceNotification, d, notification)); this.initDeviceRefresher(device); - this.eventEmitter.emit(DeviceManagerEvent.deviceConnected, device); - - this.claimDetectedDevice(deviceInfo.detectionId); - - return true; - } + this.connectedDevices.set(device.getDeviceId, device); - // Keyed by detection id so revokeDetectedDevice() (which only has that id) can drop it - private registerPendingRetry( - deviceInfo: DeviceDetectionInfo, - canonicalId: DeviceId, - closingDevice?: Promise - ): void { - this.pendingRetries.set(deviceInfo.detectionId, { deviceInfo, canonicalId, closingDevice }); + this.eventEmitter.emit(DeviceManagerEvent.deviceConnected, device); } public async onSettingsChanged(): Promise { @@ -211,27 +181,20 @@ export default class DeviceManager } } - for (const [detectionId, { deviceInfo, canonicalId, closingDevice }] of this.pendingRetries) { - if (!this.isDeviceEnabled(canonicalId)) { + for (const [detectionId, disabledDetectedDevice] of [...this.detectedDisabledDevices]) { + if (!this.isDeviceEnabled(disabledDetectedDevice.canonicalId)) { continue; } - this.pendingRetries.delete(detectionId); + this.detectedDisabledDevices.delete(detectionId); - // Make sure a device rejected by addDevice() has finished closing before re-announcing - if (undefined !== closingDevice) { - await closingDevice; - } + // Make sure a device rejected for being disabled has finished closing before re-announcing + await disabledDetectedDevice.deviceReleased; - this.announceDetectedDevice(deviceInfo); + this.announceDetectedDevice(disabledDetectedDevice.deviceDetectionInfo); } } - public claimDetectedDevice(deviceId: DeviceId): void - { - this.clearDetectedDeviceAcquireQueue(deviceId, `Device with id '${deviceId}' has been claimed by another provider`); - } - public getConnectedDevices(): AnyDevice[] { return Array.from(this.connectedDevices.values()); @@ -262,6 +225,8 @@ export default class DeviceManager public async reset(): Promise { + this.offerQueue.closeAll(new DeviceOfferRejectedError('Device manager reset')); + let closeError: unknown; for (const [, device] of this.connectedDevices) { @@ -275,26 +240,13 @@ export default class DeviceManager } } - for (const [deviceId] of this.detectedDeviceAcquireQueue) { - this.clearDetectedDeviceAcquireQueue(deviceId, 'Device manager reset'); - } - - this.pendingRetries.clear(); + this.detectedDisabledDevices.clear(); if (undefined !== closeError) { throw closeError; } } - private clearDetectedDeviceAcquireQueue(deviceId: string, reason: string): void - { - for (const entry of this.detectedDeviceAcquireQueue.get(deviceId) ?? []) { - entry.resolve({ successful: false, reason }); - } - - this.detectedDeviceAcquireQueue.delete(deviceId); - } - private initDeviceRefresher(device: AnyDevice): void { this.logger.info(`Initializing refresher for device '${device.getDeviceName}' (id: ${device.getDeviceId})`); const deviceRefreshIntervalMs = device.getRefreshInterval; @@ -322,15 +274,4 @@ export default class DeviceManager device.on(DeviceEvent.deviceDisconnected, () => deviceRefreshInterval.clear()); } - - private removeDevice(device: AnyDevice): void - { - this.connectedDevices.delete(device.getDeviceId); - this.eventEmitter.emit(DeviceManagerEvent.deviceDisconnected, device); - } - - private refreshDevice(device: AnyDevice): void - { - this.eventEmitter.emit(DeviceManagerEvent.deviceRefreshed, device); - } } diff --git a/src/device/deviceOfferRejectedError.ts b/src/device/deviceOfferRejectedError.ts new file mode 100644 index 00000000..950592e4 --- /dev/null +++ b/src/device/deviceOfferRejectedError.ts @@ -0,0 +1 @@ +export default class DeviceOfferRejectedError extends Error {} diff --git a/src/device/protocol/airotic/airoticDeviceProvider.ts b/src/device/protocol/airotic/airoticDeviceProvider.ts index ede27b87..2a483d40 100644 --- a/src/device/protocol/airotic/airoticDeviceProvider.ts +++ b/src/device/protocol/airotic/airoticDeviceProvider.ts @@ -1,4 +1,3 @@ -import EventEmitter from 'events'; import BaseError from 'modern-errors'; import DeviceManager from '../../deviceManager.js'; import AiroticDevice from './airoticDevice.js'; @@ -24,15 +23,14 @@ export default class AiroticDeviceProvider extends BleDeviceProvider { + protected override async connectBleDevice(deviceDetectionInfo: BleDeviceDetectionInfo): Promise { const transport = await promiseWithTimeout(BleUartDeviceTransport.create( deviceDetectionInfo.peripheral, AiroticDeviceProvider.UART_RX_CHAR_UUID, @@ -47,8 +45,15 @@ export default class AiroticDeviceProvider extends BleDeviceProvider { + protected override createDevice(deviceDetectionInfo: ButtplugIoDeviceDetectionInfo): Promise { const device = this.buttplugIoDeviceFactory.create( deviceDetectionInfo.detectionId, deviceDetectionInfo.buttplugClientDevice, diff --git a/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProviderFactory.ts b/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProviderFactory.ts index 00ac3d00..0a0b36e7 100644 --- a/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProviderFactory.ts +++ b/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProviderFactory.ts @@ -1,4 +1,3 @@ -import EventEmitter from 'events'; import DeviceProviderFactory from '../../provider/deviceProviderFactory.js'; import Logger from '../../../logging/Logger.js'; import ButtplugIoDeviceFactory from './buttplugIoDeviceFactory.js'; @@ -15,20 +14,16 @@ export default class ButtplugIoWebsocketDeviceProviderFactory implements DeviceP { private readonly deviceManager: DeviceManager; - private readonly eventEmitter: EventEmitter; - private readonly deviceFactory: ButtplugIoDeviceFactory; private readonly logger: Logger; public constructor( deviceManager: DeviceManager, - eventEmitter: EventEmitter, deviceFactory: ButtplugIoDeviceFactory, logger: Logger ) { this.deviceManager = deviceManager; - this.eventEmitter = eventEmitter; this.deviceFactory = deviceFactory; this.logger = logger; } @@ -37,7 +32,6 @@ export default class ButtplugIoWebsocketDeviceProviderFactory implements DeviceP { return new ButtplugIoWebsocketDeviceProvider( this.deviceManager, - this.eventEmitter, this.deviceFactory, config.address, config.autoScan, diff --git a/src/device/protocol/estim2b/estim2bSerialDeviceProvider.ts b/src/device/protocol/estim2b/estim2bSerialDeviceProvider.ts index 156028ad..e28d42bf 100644 --- a/src/device/protocol/estim2b/estim2bSerialDeviceProvider.ts +++ b/src/device/protocol/estim2b/estim2bSerialDeviceProvider.ts @@ -1,7 +1,6 @@ import { ReadlineParser } from 'serialport'; import { SerialPortStream } from '@serialport/stream'; import { BindingInterface } from '@serialport/bindings-interface'; -import EventEmitter from 'events'; import Logger from '../../../logging/Logger.js'; import SerialDeviceProvider, { SerialDeviceProviderPortOpenOptions } from '../../provider/serialDeviceProvider.js'; import EStim2bProtocol from './estim2bProtocol.js'; @@ -27,17 +26,16 @@ export default class EStim2bSerialDeviceProvider extends SerialDeviceProvider): Promise { + protected async connectSerialDevice(deviceDetectionInfo: SerialDeviceDetectionInfo, port: SerialPortStream): Promise { const parser = port.pipe(new ReadlineParser({ delimiter: '\n' })); const syncPort = new SynchronousSerialPort(deviceDetectionInfo.portInfo, parser, port, this.logger); const transport = this.transportFactory.create(syncPort, undefined, Buffer.from('\r')); diff --git a/src/device/protocol/slvCtrlPlus/slvCtrlPlusSerialDeviceProvider.ts b/src/device/protocol/slvCtrlPlus/slvCtrlPlusSerialDeviceProvider.ts index f7cb6c79..c2d18a69 100644 --- a/src/device/protocol/slvCtrlPlus/slvCtrlPlusSerialDeviceProvider.ts +++ b/src/device/protocol/slvCtrlPlus/slvCtrlPlusSerialDeviceProvider.ts @@ -3,7 +3,6 @@ import { SerialPortStream } from '@serialport/stream'; import { BindingInterface, PortInfo } from '@serialport/bindings-interface'; import SlvCtrlPlusDeviceFactory from './slvCtrlPlusDeviceFactory.js'; import SynchronousSerialPort from '../../../serial/synchronousSerialPort.js'; -import EventEmitter from 'events'; import SerialDeviceTransportFactory from '../../transport/serialDeviceTransportFactory.js'; import Logger from '../../../logging/Logger.js'; import SerialDeviceProvider, { SerialDeviceProviderPortOpenOptions } from '../../provider/serialDeviceProvider.js'; @@ -31,17 +30,16 @@ export default class SlvCtrlPlusSerialDeviceProvider extends SerialDeviceProvide deviceManager: DeviceManager, serialPortFactory: SerialPortFactory, serialPortObserver: SerialPortObserver, - eventEmitter: EventEmitter, deviceFactory: SlvCtrlPlusDeviceFactory, deviceTransportFactory: SerialDeviceTransportFactory, logger: Logger ) { - super(deviceManager, serialPortFactory, serialPortObserver, eventEmitter, logger.child({ name: SlvCtrlPlusSerialDeviceProvider.name })); + super(deviceManager, serialPortFactory, serialPortObserver, logger.child({ name: SlvCtrlPlusSerialDeviceProvider.name })); this.slvCtrlPlusDeviceFactory = deviceFactory; this.deviceTransportFactory = deviceTransportFactory; } - protected async connectSerialDevice(deviceDetectionInfo: SerialDeviceDetectionInfo, port: SerialPortStream): Promise + protected async connectSerialDevice(deviceDetectionInfo: SerialDeviceDetectionInfo, port: SerialPortStream): Promise { const parser = port.pipe(new ReadlineParser({ delimiter: SlvCtrlProtocol.EOF })); const syncPort = new SynchronousSerialPort(deviceDetectionInfo.portInfo, parser, port, this.logger); diff --git a/src/device/protocol/virtual/virtualDeviceProvider.ts b/src/device/protocol/virtual/virtualDeviceProvider.ts index 5a4f07ac..d6e7ece3 100644 --- a/src/device/protocol/virtual/virtualDeviceProvider.ts +++ b/src/device/protocol/virtual/virtualDeviceProvider.ts @@ -1,4 +1,3 @@ -import EventEmitter from 'events'; import DeviceProvider from '../../provider/deviceProvider.js'; import Logger from '../../../logging/Logger.js'; import VirtualDevice from './virtualDevice.js'; @@ -29,12 +28,11 @@ export default class VirtualDeviceProvider extends DeviceProvider | undefined> { + protected override createDevice(deviceDetectionInfo: VirtualDeviceDetectionInfo): Promise> { this.logger.info(`Virtual device detected: ${deviceDetectionInfo.knownDevice.name}`, deviceDetectionInfo.knownDevice); return this.deviceFactory.create(deviceDetectionInfo.knownDevice, VirtualDeviceProvider.providerName); diff --git a/src/device/protocol/virtual/virtualDeviceProviderFactory.ts b/src/device/protocol/virtual/virtualDeviceProviderFactory.ts index a92a2e55..398d219d 100644 --- a/src/device/protocol/virtual/virtualDeviceProviderFactory.ts +++ b/src/device/protocol/virtual/virtualDeviceProviderFactory.ts @@ -4,14 +4,11 @@ import VirtualDeviceProvider from './virtualDeviceProvider.js'; import SettingsManager from '../../../settings/settingsManager.js'; import VirtualDeviceFactory from './virtualDeviceFactory.js'; import DeviceManager from '../../deviceManager.js'; -import EventEmitterFactory from '../../../factory/eventEmitterFactory.js'; export default class VirtualDeviceProviderFactory implements DeviceProviderFactory { private readonly deviceManager: DeviceManager; - private readonly eventEmitterFactory: EventEmitterFactory; - private readonly deviceFactory: VirtualDeviceFactory; private readonly settingsManager: SettingsManager; @@ -20,13 +17,11 @@ export default class VirtualDeviceProviderFactory implements DeviceProviderFacto public constructor( deviceManager: DeviceManager, - eventEmitterFactory: EventEmitterFactory, deviceFactory: VirtualDeviceFactory, settingsManager: SettingsManager, logger: Logger ) { this.deviceManager = deviceManager; - this.eventEmitterFactory = eventEmitterFactory; this.deviceFactory = deviceFactory; this.settingsManager = settingsManager; this.logger = logger; @@ -35,7 +30,6 @@ export default class VirtualDeviceProviderFactory implements DeviceProviderFacto public create(): VirtualDeviceProvider { return new VirtualDeviceProvider( this.deviceManager, - this.eventEmitterFactory.create(), this.deviceFactory, this.settingsManager, this.logger, diff --git a/src/device/protocol/zc95/zc95SerialDeviceProvider.ts b/src/device/protocol/zc95/zc95SerialDeviceProvider.ts index 52174b91..f340a403 100644 --- a/src/device/protocol/zc95/zc95SerialDeviceProvider.ts +++ b/src/device/protocol/zc95/zc95SerialDeviceProvider.ts @@ -1,6 +1,5 @@ import { SerialPortStream } from '@serialport/stream'; import { BindingInterface } from '@serialport/bindings-interface'; -import EventEmitter from 'events'; import Logger from '../../../logging/Logger.js'; import SerialDeviceProvider, { SerialDeviceProviderPortOpenOptions } from '../../provider/serialDeviceProvider.js'; import Zc95DeviceFactory from './zc95DeviceFactory.js'; @@ -28,17 +27,16 @@ export default class Zc95SerialDeviceProvider extends SerialDeviceProvider): Promise { + protected async connectSerialDevice(deviceDetectionInfo: SerialDeviceDetectionInfo, port: SerialPortStream): Promise { const serialLogger = this.logger.child({ name: Zc95Device.name }) const parser = port.pipe(new FrameParser({ stx: Zc95Protocol.STX, etx: Zc95Protocol.ETX })); diff --git a/src/device/provider/bleDeviceProvider.ts b/src/device/provider/bleDeviceProvider.ts index e8da8c14..3df661d6 100644 --- a/src/device/provider/bleDeviceProvider.ts +++ b/src/device/provider/bleDeviceProvider.ts @@ -1,4 +1,3 @@ -import EventEmitter from 'events'; import { Peripheral } from '@stoprocent/noble'; import DeviceProvider from './deviceProvider.js'; import DeviceManager, { DeviceDetectionInfo } from '../deviceManager.js'; @@ -12,8 +11,8 @@ export default abstract class BleDeviceProvider extends { private readonly bleObserver: BleObserver; - protected constructor(deviceManager: DeviceManager, bleObserver: BleObserver, eventEmitter: EventEmitter, logger: Logger) { - super(deviceManager, eventEmitter, logger); + protected constructor(deviceManager: DeviceManager, bleObserver: BleObserver, logger: Logger) { + super(deviceManager, logger); this.bleObserver = bleObserver; } @@ -29,7 +28,7 @@ export default abstract class BleDeviceProvider extends return deviceDetectionInfo.type === 'ble'; } - protected override createDevice(deviceDetectionInfo: BleDeviceDetectionInfo): Promise { + protected override createDevice(deviceDetectionInfo: BleDeviceDetectionInfo): Promise { return this.connectBleDevice(deviceDetectionInfo); } @@ -53,5 +52,5 @@ export default abstract class BleDeviceProvider extends } } - protected abstract connectBleDevice(deviceDetectionInfo: BleDeviceDetectionInfo): Promise; + protected abstract connectBleDevice(deviceDetectionInfo: BleDeviceDetectionInfo): Promise; } diff --git a/src/device/provider/deviceProvider.ts b/src/device/provider/deviceProvider.ts index 0cb0f9f6..94bb51b1 100644 --- a/src/device/provider/deviceProvider.ts +++ b/src/device/provider/deviceProvider.ts @@ -1,10 +1,11 @@ -import EventEmitter from 'events'; import DeviceManager, { DeviceDetectionInfo, DeviceManagerEvent } from '../deviceManager.js'; +import DeviceOfferRejectedError from '../deviceOfferRejectedError.js'; import Logger from '../../logging/Logger.js'; import { asyncHandler } from '../../util/async.js'; import { logError } from '../../util/error.js'; import { AnyDevice, DeviceEvent } from '../device.js'; import { DeviceId } from '../deviceId.js'; +import BaseError from 'modern-errors'; export type AnyDeviceProvider = DeviceProvider; @@ -12,8 +13,6 @@ export default abstract class DeviceProvider = new Map(); @@ -22,9 +21,8 @@ export default abstract class DeviceProvider this.createAndRegisterDevice(deviceDetectionInfo)); - if (!acquireResult.successful) { - this.logger.debug(`Could not acquire device: ${acquireResult.reason}`); - return; - } + if (!result.successful) { + if (result.reason instanceof DeviceOfferRejectedError) { + this.logger.info(`Offer for device detection with id '${deviceDetectionInfo.detectionId}' was rejected: ${result.reason.message}`); + } else { + this.logger.info(`Offer for device detection with id '${deviceDetectionInfo.detectionId}' failed: ${BaseError.normalize(result.reason).message}`); - let device: D | undefined; - - try { - device = await this.createDevice(deviceDetectionInfo); - } catch (e: unknown) { - logError(this.logger, `Error while connecting to device '${deviceDetectionInfo.detectionId}'`, e); - await this.abortDetection(deviceDetectionInfo); - return; + // Only a real connect failure (a thrown offer) warrants provider cleanup + await this.onConnectFailed(deviceDetectionInfo); + } } + } + + private async createAndRegisterDevice(deviceDetectionInfo: DDI): Promise + { + const device = await this.createDevice(deviceDetectionInfo); - if (undefined === device || this.isStopped()) { + // Provider was stopped while the offer was in flight + if (this.isStopped()) { try { - if (undefined !== device) { - await device.close(); - } - } finally { - await this.abortDetection(deviceDetectionInfo); + await device.close(); + } catch (e: unknown) { + logError(this.logger, `Failed to close device '${device.getDeviceId}' after provider was stopped`, e); } - return; + throw new Error(`Provider was stopped while connecting device '${deviceDetectionInfo.detectionId}'`); } - device.on(DeviceEvent.deviceDisconnected, (d) => this.connectedDevices.delete(d.getDeviceId)); - - // The device manager may reject the device, e.g. because it is disabled - if (!this.deviceManager.addDevice(deviceDetectionInfo, device)) { - return; - } + device.on(DeviceEvent.deviceDisconnected, (d) => { + this.connectedDevices.delete(d.getDeviceId); + this.logger.info(`Connected devices: ${this.connectedDevices.size}`); + }); this.connectedDevices.set(device.getDeviceId, device); - this.logger.info(`Connected devices: ${this.connectedDevices.size}`); - } - private async abortDetection(deviceDetectionInfo: DDI): Promise { - try { - await this.onConnectFailed(deviceDetectionInfo); - } finally { - this.deviceManager.releaseDetectedDevice(deviceDetectionInfo.detectionId); - } + return device; } protected abstract canHandleDeviceDetectionInfo(deviceDetectionInfo: DeviceDetectionInfo): deviceDetectionInfo is DDI; - protected abstract createDevice(deviceDetectionInfo: DDI): Promise; + protected abstract createDevice(deviceDetectionInfo: DDI): Promise; // eslint-disable-next-line @typescript-eslint/no-unused-vars protected async onConnectFailed(deviceDetectionInfo: DDI): Promise { diff --git a/src/device/provider/deviceProviderManager.ts b/src/device/provider/deviceProviderManager.ts index be6617e6..0bd91979 100644 --- a/src/device/provider/deviceProviderManager.ts +++ b/src/device/provider/deviceProviderManager.ts @@ -1,4 +1,4 @@ -import { SequentialTaskQueue } from 'sequential-task-queue'; +import { SequentialTaskQueue } from '@timesplinter/sequential-task-queue'; import Settings from '../../settings/settings.js'; import DeviceSource from '../../settings/deviceSource.js'; import DeviceProviderFactory from './deviceProviderFactory.js'; diff --git a/src/device/provider/serialDeviceProvider.ts b/src/device/provider/serialDeviceProvider.ts index f6954b7d..f6582380 100644 --- a/src/device/provider/serialDeviceProvider.ts +++ b/src/device/provider/serialDeviceProvider.ts @@ -1,5 +1,4 @@ import DeviceProvider from './deviceProvider.js'; -import EventEmitter from 'events'; import Logger from '../../logging/Logger.js'; import { BindingInterface, PortInfo } from '@serialport/bindings-interface'; import { SerialPortOpenOptions } from 'serialport'; @@ -24,10 +23,9 @@ export default abstract class SerialDeviceProvider { + protected override async createDevice(deviceDetectionInfo: SerialDeviceDetectionInfo): Promise { const portInfo = deviceDetectionInfo.portInfo; this.logger.info(`Connection attempt for serial device '${portInfo.path}' (s/n: ${portInfo.serialNumber})`); @@ -56,9 +54,6 @@ export default abstract class SerialDeviceProvider((resolve, reject) => { port.open(err => err ? reject(err) : resolve()); @@ -66,32 +61,28 @@ export default abstract class SerialDeviceProvider((resolve, reject) => { + port.close(err => err ? reject(err) : resolve()); + }); } catch (closeError: unknown) { - logError(this.logger, `Failed to close partially registered serial device '${portInfo.path}'`, closeError); + logError(this.logger, `Failed to close serial port '${portInfo.path}' after a failed connection attempt`, closeError); } } + const error = BaseError.normalize(e); - attemptFailureReason = error.message; - } + this.logger.info(`Could not connect to serial device '${portInfo.path}': ${error.message}`); - if (undefined === device) { - if (port.isOpen) { - await new Promise((resolve, reject) => { - port.close(err => err ? reject(err) : resolve()); - }); - } - this.logger.info(`Could not connect to serial device '${portInfo.path}': ${attemptFailureReason}`); - } else { - this.logger.info(`Successfully connected to serial device '${portInfo.path}'`); - this.logger.debug(`Assigned device id: ${device.getDeviceId} (${portInfo.path})`); + throw e; } - - return device; } // eslint-disable-next-line @typescript-eslint/no-unused-vars @@ -99,7 +90,7 @@ export default abstract class SerialDeviceProvider): Promise; + protected abstract connectSerialDevice(deviceDetectionInfo: SerialDeviceDetectionInfo, port: SerialPortStream): Promise; protected abstract getSerialDeviceProviderPortOpenOptions(portInfo: PortInfo): SerialDeviceProviderPortOpenOptions; } diff --git a/src/device/updater/bufferedDeviceUpdater.ts b/src/device/updater/bufferedDeviceUpdater.ts index 1f21669c..b58b97bf 100644 --- a/src/device/updater/bufferedDeviceUpdater.ts +++ b/src/device/updater/bufferedDeviceUpdater.ts @@ -1,6 +1,6 @@ import { AnyDevice, DeviceData } from '../device.js'; import DeviceUpdaterInterface from './deviceUpdaterInterface.js'; -import { SequentialTaskQueue } from 'sequential-task-queue'; +import { SequentialTaskQueue } from '@timesplinter/sequential-task-queue'; export default class BufferedDeviceUpdater implements DeviceUpdaterInterface { diff --git a/src/serial/synchronousSerialPort.ts b/src/serial/synchronousSerialPort.ts index 39616b81..4db0ad65 100644 --- a/src/serial/synchronousSerialPort.ts +++ b/src/serial/synchronousSerialPort.ts @@ -1,5 +1,5 @@ import { Readable, Writable } from 'stream'; -import { cancellationTokenReasons, SequentialTaskQueue, TaskOptions } from 'sequential-task-queue'; +import { cancellationTokenReasons, SequentialTaskQueue, TaskOptions } from '@timesplinter/sequential-task-queue'; import { PortInfo } from '@serialport/bindings-interface'; import Logger from '../logging/Logger.js'; import { asyncHandler } from '../util/async.js'; diff --git a/src/serviceProvider/deviceServiceProvider.ts b/src/serviceProvider/deviceServiceProvider.ts index 6317f7ad..43809239 100644 --- a/src/serviceProvider/deviceServiceProvider.ts +++ b/src/serviceProvider/deviceServiceProvider.ts @@ -58,7 +58,6 @@ export default class DeviceServiceProvider implements ServiceProvider new ButtplugIoWebsocketDeviceProviderFactory( container.get('device.manager'), - container.get('factory.eventEmitter').create(), container.get('device.serial.factory.buttplugIo'), container.get('logger.default'), ) @@ -138,7 +136,6 @@ export default class DeviceServiceProvider implements ServiceProvider new VirtualDeviceProviderFactory( container.get('device.manager'), - container.get('factory.eventEmitter'), container.get('device.virtual.factory'), container.get('settings.manager'), container.get('logger.default'), @@ -223,7 +220,6 @@ export default class DeviceServiceProvider implements ServiceProvider { + let mockedLogger: ReturnType>; + const deviceId = DeviceId.create('device-1'); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; + + beforeEach(() => { + mockedLogger = mock(); + mockedLogger.child.mockReturnValue(mockedLogger); + }); + + describe('has', () => { + it('is false before any offer and true while one is pending', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + expect(queue.has(deviceId)).toBe(false); + + const pendingPromise = queue.offer(deviceInfo, () => new Promise(() => {})); + + await vi.waitFor(() => expect(queue.has(deviceId)).toBe(true)); + + queue.closeAll(new DeviceOfferRejectedError('test cleanup')); + await pendingPromise; + }); + }); + + describe('offer', () => { + it('lazily opens a queue and runs the offer even without any prior activity for the detection id', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + + const result = await queue.offer(deviceInfo, () => Promise.resolve(device)); + + expect(result).toStrictEqual({ successful: true, device }); + }); + + it('does not run a second offer while the first is still pending', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + let resolveFirstOffer!: (device: AnyDevice) => void; + const firstOfferPromise = new Promise((resolve) => { resolveFirstOffer = resolve; }); + const secondOfferFn = vi.fn(() => Promise.resolve(new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()))); + + const firstResultPromise = queue.offer(deviceInfo, () => firstOfferPromise); + queue.offer(deviceInfo, secondOfferFn); + + expect(secondOfferFn).not.toHaveBeenCalled(); + + resolveFirstOffer(new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter())); + await firstResultPromise; + }); + + it('hands off to the next queued offer when the first one throws', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const offerError = new Error('connection failed'); + + const firstResultPromise = queue.offer(deviceInfo, () => Promise.reject(offerError)); + const secondResultPromise = queue.offer(deviceInfo, () => Promise.resolve(device)); + + const [firstResult, secondResult] = await Promise.all([firstResultPromise, secondResultPromise]); + + expect(firstResult).toStrictEqual({ successful: false, reason: offerError }); + expect(secondResult).toStrictEqual({ successful: true, device }); + }); + + it('hands off to the next queued offer when the first device is rejected', async () => { + const rejection = new DeviceOfferRejectedError('rejected'); + const secondDevice = new TestDevice(DeviceId.create('device-1-second'), 'Foo', new Date(), false, new EventEmitter()); + + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + const firstResultPromise = queue.offer(deviceInfo, () => Promise.resolve(rejection)); + const secondResultPromise = queue.offer(deviceInfo, () => Promise.resolve(secondDevice)); + + const [firstResult, secondResult] = await Promise.all([firstResultPromise, secondResultPromise]); + + expect(firstResult).toStrictEqual({ successful: false, reason: rejection }); + expect(secondResult).toStrictEqual({ successful: true, device: secondDevice }); + }); + + it('lazily reopens for a new offer after the only queued offer is rejected', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + await queue.offer(deviceInfo, () => Promise.resolve(new DeviceOfferRejectedError('rejected'))); + + expect(queue.has(deviceId)).toBe(false); + + const pendingPromise = queue.offer(deviceInfo, () => new Promise(() => {})); + expect(queue.has(deviceId)).toBe(true); + + queue.closeAll(new DeviceOfferRejectedError('test cleanup')); + await pendingPromise; + }); + + it('rejects other queued offers with DeviceOfferRejectedError once a device is claimed', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + let resolveFirstOffer!: (device: AnyDevice) => void; + const firstOfferPromise = new Promise((resolve) => { resolveFirstOffer = resolve; }); + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + + const firstResultPromise = queue.offer(deviceInfo, () => firstOfferPromise); + const secondResultPromise = queue.offer(deviceInfo, () => Promise.reject(new Error('should never run'))); + + resolveFirstOffer(device); + + const [firstResult, secondResult] = await Promise.all([firstResultPromise, secondResultPromise]); + + expect(firstResult).toStrictEqual({ successful: true, device }); + expect(secondResult.successful).toBe(false); + expect(!secondResult.successful && secondResult.reason).toBeInstanceOf(DeviceOfferRejectedError); + }); + }); + + describe('closeAll', () => { + // closeAll() replaces clear()'s narrower "just this one id" purpose - these exercise the + // same underlying close()/cancel() mechanics, scoped to a single queue at a time. + + it('resolves a pending offer with failure', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + // First offer never settles on its own, so it's still holding the queue when closed + const pendingPromise = queue.offer(deviceInfo, () => new Promise(() => {})); + + queue.closeAll(new DeviceOfferRejectedError('revoked')); + + const result = await pendingPromise; + expect(result.successful).toBe(false); + expect(!result.successful && result.reason).toBeInstanceOf(DeviceOfferRejectedError); + }); + + it('closes a device whose offer resolves after the queue was cleared, without accepting it', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + let resolveOffer!: (device: AnyDevice) => void; + const offerPromise = new Promise((resolve) => { resolveOffer = resolve; }); + let offerStarted = false; + const resultPromise = queue.offer(deviceInfo, () => { + offerStarted = true; + return offerPromise; + }); + + // Wait for the offer to actually start running (SequentialTaskQueue starts tasks via + // its scheduler, not synchronously) before clearing, to genuinely simulate a revoke + // while the offer is in flight rather than while it's still merely queued. + await vi.waitFor(() => expect(offerStarted).toBe(true)); + + // Device physically disappears while the offer is still in flight. + queue.closeAll(new DeviceOfferRejectedError('revoked')); + + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const closeSpy = vi.spyOn(device, 'close'); + + // The offer only settles now, after the queue was already cleared - the caller + // already got a rejected result above, so this device must never be accepted. + resolveOffer(device); + + const result = await resultPromise; + expect(result.successful).toBe(false); + expect(!result.successful && result.reason).toBeInstanceOf(DeviceOfferRejectedError); + + await vi.waitFor(() => expect(closeSpy).toHaveBeenCalled()); + }); + + it('does not corrupt a fresh, still-pending queue when a stale offer fails after a clear', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + let rejectStaleOffer!: (reason: unknown) => void; + const staleOfferPromise = new Promise((_resolve, reject) => { rejectStaleOffer = reject; }); + let staleOfferStarted = false; + const staleResultPromise = queue.offer(deviceInfo, () => { + staleOfferStarted = true; + return staleOfferPromise; + }); + + // Wait for the stale offer to actually start running before clearing, to genuinely + // simulate a revoke while it's in flight rather than while it's still merely queued. + await vi.waitFor(() => expect(staleOfferStarted).toBe(true)); + + // Device physically disappears while the stale offer is still in flight. + queue.closeAll(new DeviceOfferRejectedError('revoked')); + + // Re-detected under the same detection id - offer() lazily opens a fresh queue, with + // its own still-pending offer. + let resolveFreshOffer!: (device: AnyDevice) => void; + const freshOfferPromise = new Promise((resolve) => { resolveFreshOffer = resolve; }); + const freshResultPromise = queue.offer(deviceInfo, () => freshOfferPromise); + + // The stale offer only fails now, well after it was cleared and superseded - while + // the fresh offer is still pending. Without the staleness guard, this would + // incorrectly wipe the still-valid, currently in-flight fresh queue. + rejectStaleOffer(new Error('stale offer failed')); + await staleResultPromise; + + // A second offer for the same detection id right now must still be queued up behind + // the fresh offer, not told the device is unavailable (which is what would happen if + // the stale processing had wrongly wiped the still-valid queue). + const secondResultPromise = queue.offer(deviceInfo, () => Promise.reject(new Error('should never run'))); + + const freshDevice = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + resolveFreshOffer(freshDevice); + + const [freshResult, secondResult] = await Promise.all([freshResultPromise, secondResultPromise]); + + expect(freshResult).toStrictEqual({ successful: true, device: freshDevice }); + expect(secondResult.successful).toBe(false); + // Specifically "claimed by another provider" (queued behind the still-valid fresh + // queue) - not "not available anymore for offering", which would mean the queue was + // wrongly wiped by the stale processing. + expect(!secondResult.successful && (secondResult.reason as Error).message).toContain('claimed by another provider'); + }); + + it('closes every open queue', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + const otherDeviceId = DeviceId.create('device-2'); + + const firstResultPromise = queue.offer(deviceInfo, () => new Promise(() => {})); + const secondResultPromise = queue.offer({ type: 'test', detectionId: otherDeviceId }, () => new Promise(() => {})); + + queue.closeAll(new DeviceOfferRejectedError('reset')); + + const [firstResult, secondResult] = await Promise.all([firstResultPromise, secondResultPromise]); + + expect(firstResult.successful).toBe(false); + expect(secondResult.successful).toBe(false); + expect(queue.has(deviceId)).toBe(false); + expect(queue.has(otherDeviceId)).toBe(false); + }); + }); + + describe('revoke / dropIfRevoked', () => { + it('blocks a subsequent offer even when nothing was ever offered before the revoke', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + queue.revoke(deviceId, new DeviceOfferRejectedError('device disappeared')); + + const result = await queue.offer(deviceInfo, () => Promise.reject(new Error('should never run'))); + + expect(result.successful).toBe(false); + expect(!result.successful && result.reason).toBeInstanceOf(DeviceOfferRejectedError); + }); + + it('cancels an in-flight offer and keeps rejecting further offers after the revoke', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + let offerStarted = false; + const resultPromise = queue.offer(deviceInfo, () => { + offerStarted = true; + return new Promise(() => {}); + }); + + await vi.waitFor(() => expect(offerStarted).toBe(true)); + + queue.revoke(deviceId, new DeviceOfferRejectedError('device disappeared')); + + const result = await resultPromise; + expect(result.successful).toBe(false); + + // The tombstone must survive the drain triggered by the revoke's own cancellation - + // a late offer arriving right after must still see it and reject itself, instead of + // unknowingly reopening a queue for a device that's already confirmed gone. + const lateResult = await queue.offer(deviceInfo, () => Promise.reject(new Error('should never run'))); + + expect(lateResult.successful).toBe(false); + expect(!lateResult.successful && lateResult.reason).toBeInstanceOf(DeviceOfferRejectedError); + }); + + it('allows a fresh offer to succeed again once dropIfRevoked acknowledges a genuine redetection', async () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + queue.revoke(deviceId, new DeviceOfferRejectedError('device disappeared')); + queue.dropIfRevoked(deviceId); + + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const result = await queue.offer(deviceInfo, () => Promise.resolve(device)); + + expect(result).toStrictEqual({ successful: true, device }); + }); + + it('dropIfRevoked is a no-op when there is nothing to drop', () => { + const queue = new DetectedDeviceOfferQueue(mockedLogger); + + expect(() => queue.dropIfRevoked(deviceId)).not.toThrow(); + expect(queue.has(deviceId)).toBe(false); + }); + }); +}); diff --git a/tests/unit/device/deviceManager.spec.ts b/tests/unit/device/deviceManager.spec.ts index 77901675..3406c289 100644 --- a/tests/unit/device/deviceManager.spec.ts +++ b/tests/unit/device/deviceManager.spec.ts @@ -1,8 +1,9 @@ -import {describe, it, expect, beforeEach} from "vitest"; +import {describe, it, expect, beforeEach, vi} from "vitest"; import {mock,mockClear} from "vitest-mock-extended"; import DeviceManager, { DeviceManagerEvent, DeviceDetectionInfo } from "../../../src/device/deviceManager.js"; +import DeviceOfferRejectedError from "../../../src/device/deviceOfferRejectedError.js"; import {EventEmitter} from "events"; -import Device from "../../../src/device/device.js"; +import Device, { AnyDevice } from "../../../src/device/device.js"; import TestDevice from "./testDevice.js"; import Logger from "../../../src/logging/Logger.js"; import { DeviceId } from "../../../src/device/deviceId.js"; @@ -15,9 +16,39 @@ describe('deviceManager', () => { // device as enabled - the desired default for tests unrelated to the enable/disable feature. const mockedSettingsManager = mock(); + // Configures a mocked EventEmitter to synchronously react to deviceDetected the way a real + // DeviceProvider does (see DeviceProvider.handleDeviceDetection(), which calls offerDevice() + // synchronously within its own emit() dispatch) - announceDetectedDevice()'s post-emit + // offerQueue.has() check depends on this happening synchronously, which a bare mock doesn't + // do on its own. + const reactToDetection = (mockedEventEmitter: ReturnType>, reaction: () => void) => { + mockedEventEmitter.emit.mockImplementation((event: string | symbol) => { + if (event === DeviceManagerEvent.deviceDetected) { + reaction(); + } + return true; + }); + }; + + // Announces the device and immediately offers it for connection, simulating a provider that + // synchronously reacts to the announcement - the only way to get a device registered through + // the public API now that addDevice() is private. + const connectDevice = (manager: DeviceManager, mockedEventEmitter: ReturnType>, deviceInfo: DeviceDetectionInfo, device: AnyDevice) => { + let offerPromise!: ReturnType; + + reactToDetection(mockedEventEmitter, () => { + offerPromise = manager.offerDevice(deviceInfo, () => Promise.resolve(device)); + }); + + manager.announceDetectedDevice(deviceInfo); + + return offerPromise; + }; + it('it adds device to managed devices and emits an event', async () => { const mockedDeviceManagerEventEmitter = mock(); + mockedDeviceManagerEventEmitter.emit.mockReturnValue(true); const mockedLogger = mock(); mockedLogger.child.mockReturnValue(mockedLogger); @@ -26,18 +57,19 @@ describe('deviceManager', () => { const deviceId = DeviceId.create('test-device-id'); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; // New device connected expect(deviceManager.getConnectedDevices().length).toBe(0); - deviceManager.addDevice({ type: 'test', detectionId: deviceId }, device); + mockClear(mockedDeviceManagerEventEmitter); // drop the constructor-time noise, if any + await connectDevice(deviceManager, mockedDeviceManagerEventEmitter, deviceInfo, device); let actualDevices = deviceManager.getConnectedDevices(); expect(actualDevices.length).toBe(1); expect(actualDevices[0]).toBe(device); - expect(mockedDeviceManagerEventEmitter.emit).toBeCalledTimes(1); expect(mockedDeviceManagerEventEmitter.emit).toBeCalledWith(DeviceManagerEvent.deviceConnected, device); expect(mockedLogger.child).toBeCalledWith({ name: DeviceManager.name }); }); @@ -48,25 +80,25 @@ describe('deviceManager', () => { const connectedDevices = new Map(); const deviceId = DeviceId.create('test-device-id'); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; const mockedDeviceManagerEventEmitter = mock(); + mockedDeviceManagerEventEmitter.emit.mockReturnValue(true); const mockedLogger = mock(); mockedLogger.child.mockReturnValue(mockedLogger); const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, connectedDevices, mockedSettingsManager, mockedLogger); - deviceManager.addDevice({ type: 'test', detectionId: deviceId }, device); + await connectDevice(deviceManager, mockedDeviceManagerEventEmitter, deviceInfo, device); + mockClear(mockedDeviceManagerEventEmitter); // Connected device refreshed await device.refresh(); - expect(mockedDeviceManagerEventEmitter.emit).toBeCalledTimes(2); - expect(mockedDeviceManagerEventEmitter.emit).toHaveBeenNthCalledWith(1, DeviceManagerEvent.deviceConnected, device); - expect(mockedDeviceManagerEventEmitter.emit).toHaveBeenNthCalledWith(2, DeviceManagerEvent.deviceRefreshed, device); + expect(mockedDeviceManagerEventEmitter.emit).toBeCalledTimes(1); + expect(mockedDeviceManagerEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceRefreshed, device); expect(mockedLogger.child).toBeCalledWith({ name: DeviceManager.name }); - - mockClear(mockedDeviceManagerEventEmitter); }); it('it emits an event on device update', async () => { @@ -74,24 +106,26 @@ describe('deviceManager', () => { const connectedDevices = new Map(); const deviceId = DeviceId.create('test-device-id'); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; const mockedDeviceManagerEventEmitter = mock(); + mockedDeviceManagerEventEmitter.emit.mockReturnValue(true); const mockedLogger = mock(); mockedLogger.child.mockReturnValue(mockedLogger); const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, connectedDevices, mockedSettingsManager, mockedLogger); - deviceManager.addDevice({ type: 'test', detectionId: deviceId }, device); + await connectDevice(deviceManager, mockedDeviceManagerEventEmitter, deviceInfo, device); + mockClear(mockedDeviceManagerEventEmitter); // Connected device closed await device.close(); expect(deviceManager.getConnectedDevices().length).toBe(0); - expect(mockedDeviceManagerEventEmitter.emit).toBeCalledTimes(2); - expect(mockedDeviceManagerEventEmitter.emit).toHaveBeenNthCalledWith(1, DeviceManagerEvent.deviceConnected, device); - expect(mockedDeviceManagerEventEmitter.emit).toHaveBeenNthCalledWith(2, DeviceManagerEvent.deviceDisconnected, device); + expect(mockedDeviceManagerEventEmitter.emit).toBeCalledTimes(1); + expect(mockedDeviceManagerEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDisconnected, device); expect(mockedLogger.child).toBeCalledWith({ name: DeviceManager.name }); }); @@ -140,10 +174,16 @@ describe('deviceManager', () => { expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); }); - it('does not re-announce a device already in the acquire queue', () => { - mockedEventEmitter.emit.mockReturnValue(true); + it('does not re-announce a device while an offer for it is still in flight', () => { const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); + // Offer never settles on its own, so the queue is still legitimately open (an offer + // is genuinely in progress) when the second announce comes in - that's what's under + // test here, distinct from the "nothing ever offered" behavior below. + reactToDetection(mockedEventEmitter, () => { + void manager.offerDevice(deviceInfo, () => new Promise(() => {})); + }); + manager.announceDetectedDevice(deviceInfo); manager.announceDetectedDevice(deviceInfo); @@ -151,6 +191,26 @@ describe('deviceManager', () => { expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); }); + it('allows re-announcing a device after it was revoked (tombstone must not permanently block it)', () => { + const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); + + mockedEventEmitter.emit.mockReturnValue(true); + manager.announceDetectedDevice(deviceInfo); + + // Device physically disappears - revoke() leaves a closed tombstone behind (not a + // plain delete) so a late offer arriving after this point still rejects itself. + manager.revokeDetectedDevice(deviceInfo); + + mockClear(mockedEventEmitter); + mockedEventEmitter.emit.mockReturnValue(true); + + // Genuine redetection (e.g. replugged) must not be blocked by that leftover + // tombstone - dropIfRevoked() has to run before the has() reentrancy guard sees it. + manager.announceDetectedDevice(deviceInfo); + + expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); + }); + it('does not emit event when device is already connected', () => { const connectedDevices = new Map([[deviceId, mock()]]); const manager = new DeviceManager(mockedEventEmitter, connectedDevices, mockedSettingsManager, mockedLogger); @@ -160,17 +220,43 @@ describe('deviceManager', () => { expect(mockedEventEmitter.emit).not.toHaveBeenCalled(); }); - it('removes device from queue when no listeners respond to deviceDetected', async () => { + it('still allows a later offer to succeed on its own after no listeners responded to deviceDetected', async () => { mockedEventEmitter.emit.mockReturnValue(false); const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); - const result = await manager.acquireDetectedDevice(deviceId); - expect(result.successful).toBe(false); + // announceDetectedDevice() no longer opens/reserves anything proactively - offer() + // lazily opens its own queue, so a provider calling offerDevice() later succeeds + // rather than being told the device is "not available anymore". + const result = await manager.offerDevice(deviceInfo, () => Promise.resolve(new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()))); + expect(result.successful).toBe(true); }); - it('does not emit deviceDetected for a device belonging to a disabled known device', () => { + it('discards the queue and allows re-announcing when listeners exist but none of them offer a device', () => { + // Simulates a subscribed provider whose canHandleDeviceDetectionInfo() declines this + // detection's type, so it never calls offerDevice() - hadListeners is true, but + // nothing ever offers. + mockedEventEmitter.emit.mockReturnValue(true); + const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); + + manager.announceDetectedDevice(deviceInfo); + + mockClear(mockedEventEmitter); + mockedEventEmitter.emit.mockReturnValue(true); + + manager.announceDetectedDevice(deviceInfo); + + expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); + }); + + it('still emits deviceDetected even when the detection id matches a disabled known device', () => { + // Detection id is preliminary/raw - e.g. for serial ports, multiple protocol + // providers share the same detectionId but each computes its own distinct canonical + // id via handshake. Gating here on detectionId's own enabled state would incorrectly + // block every provider, including ones whose real canonical id isn't disabled at + // all. The disabled check that actually matters happens per-canonical-id in + // addDevice(), once a provider has connected and learned the final id. mockedEventEmitter.emit.mockReturnValue(true); const settings = new Settings(); @@ -182,11 +268,11 @@ describe('deviceManager', () => { manager.announceDetectedDevice(deviceInfo); - expect(mockedEventEmitter.emit).not.toHaveBeenCalled(); + expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); }); }); - describe('acquireDetectedDevice', () => { + describe('offerDevice', () => { let mockedLogger: ReturnType>; let mockedEventEmitter: ReturnType>; const deviceId = DeviceId.create('device-2'); @@ -199,64 +285,33 @@ describe('deviceManager', () => { mockedEventEmitter.emit.mockReturnValue(true); }); - it('returns failure when device is not in the detect queue', async () => { + it('runs the first offer immediately and adds the device on success', async () => { const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); - const result = await manager.acquireDetectedDevice(deviceId); + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const result = await connectDevice(manager, mockedEventEmitter, deviceInfo, device); - expect(result.successful).toBe(false); + expect(result).toStrictEqual({ successful: true, device }); + expect(manager.getConnectedDevices()).toContain(device); }); - it('resolves immediately with success for the first caller', async () => { + it('clears the queue and re-allows announcing after the only offer fails', async () => { const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); - manager.announceDetectedDevice(deviceInfo); - const result = await manager.acquireDetectedDevice(deviceId); + let resultPromise!: ReturnType; + reactToDetection(mockedEventEmitter, () => { + resultPromise = manager.offerDevice(deviceInfo, () => Promise.reject(new Error('connect failed'))); + }); - expect(result).toStrictEqual({ successful: true }); - }); - - it('queues the second caller until the first releases', async () => { - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); + await resultPromise; - await manager.acquireDetectedDevice(deviceId); - const secondCallerPromise = manager.acquireDetectedDevice(deviceId); - manager.releaseDetectedDevice(deviceId); - - const result = await secondCallerPromise; - expect(result).toStrictEqual({ successful: true }); - }); - }); - - describe('releaseDetectedDevice', () => { - let mockedLogger: ReturnType>; - let mockedEventEmitter: ReturnType>; - const deviceId = DeviceId.create('device-3'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; - - beforeEach(() => { - mockedLogger = mock(); - mockedLogger.child.mockReturnValue(mockedLogger); - mockedEventEmitter = mock(); + mockClear(mockedEventEmitter); mockedEventEmitter.emit.mockReturnValue(true); - }); - it('is a no-op when device is not in the acquire queue', () => { - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); - - expect(() => manager.releaseDetectedDevice(DeviceId.create('unknown'))).not.toThrow(); - }); - - it('removes device from queue after the only waiter releases', async () => { - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); - await manager.acquireDetectedDevice(deviceId); - manager.releaseDetectedDevice(deviceId); - - const result = await manager.acquireDetectedDevice(deviceId); - expect(result.successful).toBe(false); + expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); }); }); @@ -273,18 +328,6 @@ describe('deviceManager', () => { mockedEventEmitter.emit.mockReturnValue(true); }); - it('resolves a pending second caller with failure', async () => { - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); - manager.announceDetectedDevice(deviceInfo); - await manager.acquireDetectedDevice(deviceId); // first caller holds - const pendingPromise = manager.acquireDetectedDevice(deviceId); // second waits - - manager.revokeDetectedDevice(deviceInfo); - - const result = await pendingPromise; - expect(result.successful).toBe(false); - }); - it('drops a disabled device from pending retry so it is not re-announced after re-enabling', async () => { const settings = new Settings(); settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, false)); @@ -294,9 +337,14 @@ describe('deviceManager', () => { const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); - // Announced while disabled -> parked in pending retry, no deviceDetected emitted. - manager.announceDetectedDevice(deviceInfo); - expect(mockedEventEmitter.emit).not.toHaveBeenCalled(); + // Announced (detection id happens to match the disabled known device) - the provider + // still gets a chance to offer it, but addDevice() rejects it once connected since + // it's disabled, parking it in pending retry. + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const result = await connectDevice(manager, mockedEventEmitter, deviceInfo, device); + expect(result.successful).toBe(false); + + mockClear(mockedEventEmitter); // Device physically disappears while still disabled. manager.revokeDetectedDevice(deviceInfo); @@ -307,31 +355,54 @@ describe('deviceManager', () => { expect(mockedEventEmitter.emit).not.toHaveBeenCalled(); }); - }); - describe('claimDetectedDevice', () => { - let mockedLogger: ReturnType>; - let mockedEventEmitter: ReturnType>; - const deviceId = DeviceId.create('device-5'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; + it('does not park a device for retry if it was revoked while the offer was still connecting', async () => { + const settings = new Settings(); + settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, false)); - beforeEach(() => { - mockedLogger = mock(); - mockedLogger.child.mockReturnValue(mockedLogger); - mockedEventEmitter = mock(); - mockedEventEmitter.emit.mockReturnValue(true); - }); + const settingsManager = mock(); + settingsManager.getSettings.mockReturnValue(settings); + + const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); + + let resolveOffer!: (device: AnyDevice) => void; + const offerPromise = new Promise((resolve) => { resolveOffer = resolve; }); + let offerStarted = false; + let resultPromise!: ReturnType; + + reactToDetection(mockedEventEmitter, () => { + resultPromise = manager.offerDevice(deviceInfo, () => { + offerStarted = true; + return offerPromise; + }); + }); - it('resolves a pending caller with failure', async () => { - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); - await manager.acquireDetectedDevice(deviceId); // first caller holds - const pendingPromise = manager.acquireDetectedDevice(deviceId); // second waits - manager.claimDetectedDevice(deviceId); + // Wait for the connect attempt to actually start before revoking, to genuinely + // simulate a revoke while it's in flight rather than while it's still merely queued. + await vi.waitFor(() => expect(offerStarted).toBe(true)); + + // Device physically disappears while the (disabled) device is still connecting. + manager.revokeDetectedDevice(deviceInfo); - const result = await pendingPromise; + // The connect attempt only succeeds now, after the revoke already settled the caller. + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + resolveOffer(device); + + const result = await resultPromise; expect(result.successful).toBe(false); + + mockClear(mockedEventEmitter); + mockedEventEmitter.emit.mockReturnValue(true); + + // Re-enabling it must NOT resurrect the gone device - it should never have been + // parked for retry in the first place, since it was already known to be gone by the + // time it "connected". + settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); + await manager.onSettingsChanged(); + + expect(mockedEventEmitter.emit).not.toHaveBeenCalled(); }); }); @@ -375,48 +446,60 @@ describe('deviceManager', () => { }); }); - describe('addDevice - disabled devices', () => { + describe('offerDevice - disabled devices', () => { let mockedLogger: ReturnType>; + let mockedEventEmitter: ReturnType>; beforeEach(() => { mockedLogger = mock(); mockedLogger.child.mockReturnValue(mockedLogger); + mockedEventEmitter = mock(); + mockedEventEmitter.emit.mockReturnValue(true); }); - it('does not register a device belonging to a disabled known device and closes it', () => { - const deviceId = DeviceId.create('disabled-device'); + it('does not register a device whose canonical id belongs to a disabled known device, and closes it', async () => { + // Detection id is unknown/enabled so announce() lets it through and the offer runs - + // the device only turns out to be disabled once its canonical id is learned, e.g. + // during a handshake. This is the only way to reach addDevice()'s own disabled-check + // through the public API now that it's private. + const detectionId = DeviceId.create('disabled-device-detection'); + const canonicalId = DeviceId.create('disabled-device-canonical'); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId }; + const settings = new Settings(); - settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, false)); + settings.addKnownDevice(new KnownDevice(canonicalId, 'Foo', 'test', 'test', {}, false)); const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); const connectedDevices = new Map(); - const manager = new DeviceManager(mock(), connectedDevices, settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, connectedDevices, settingsManager, mockedLogger); - const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const device = new TestDevice(canonicalId, 'Foo', new Date(), false, new EventEmitter()); - const added = manager.addDevice({ type: 'test', detectionId: deviceId }, device); + const result = await connectDevice(manager, mockedEventEmitter, deviceInfo, device); - expect(added).toBe(false); + expect(result.successful).toBe(false); + expect(!result.successful && result.reason).toBeInstanceOf(DeviceOfferRejectedError); expect(manager.getConnectedDevices()).toHaveLength(0); }); - it('registers a device belonging to an enabled known device', () => { + it('registers a device belonging to an enabled known device', async () => { const deviceId = DeviceId.create('enabled-device'); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; const settings = new Settings(); settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); - const manager = new DeviceManager(mock(), new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); - const added = manager.addDevice({ type: 'test', detectionId: deviceId }, device); + const result = await connectDevice(manager, mockedEventEmitter, deviceInfo, device); - expect(added).toBe(true); + expect(result.successful).toBe(true); expect(manager.getConnectedDevices()).toHaveLength(1); }); }); @@ -431,6 +514,7 @@ describe('deviceManager', () => { it('closes connected devices whose known device has been disabled', async () => { const deviceId = DeviceId.create('device-to-disable'); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; const enabledSettings = new Settings(); enabledSettings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); @@ -438,10 +522,12 @@ describe('deviceManager', () => { settingsManager.getSettings.mockReturnValue(enabledSettings); const connectedDevices = new Map(); - const manager = new DeviceManager(mock(), connectedDevices, settingsManager, mockedLogger); + const mockedEventEmitter = mock(); + mockedEventEmitter.emit.mockReturnValue(true); + const manager = new DeviceManager(mockedEventEmitter, connectedDevices, settingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); - manager.addDevice({ type: 'test', detectionId: deviceId }, device); + await connectDevice(manager, mockedEventEmitter, deviceInfo, device); expect(manager.getConnectedDevices()).toHaveLength(1); const disabledSettings = new Settings(); @@ -455,6 +541,7 @@ describe('deviceManager', () => { it('leaves devices belonging to still-enabled known devices connected', async () => { const deviceId = DeviceId.create('device-still-enabled'); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; const settings = new Settings(); settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); @@ -462,44 +549,19 @@ describe('deviceManager', () => { settingsManager.getSettings.mockReturnValue(settings); const connectedDevices = new Map(); - const manager = new DeviceManager(mock(), connectedDevices, settingsManager, mockedLogger); - - const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); - manager.addDevice({ type: 'test', detectionId: deviceId }, device); - - await manager.onSettingsChanged(); - - expect(manager.getConnectedDevices()).toHaveLength(1); - }); - - it('re-announces a device rejected by announceDetectedDevice once its known device gets re-enabled', async () => { - const deviceId = DeviceId.create('device-pending-1'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; - - const disabledSettings = new Settings(); - disabledSettings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, false)); - - const settingsManager = mock(); - settingsManager.getSettings.mockReturnValue(disabledSettings); - const mockedEventEmitter = mock(); mockedEventEmitter.emit.mockReturnValue(true); + const manager = new DeviceManager(mockedEventEmitter, connectedDevices, settingsManager, mockedLogger); - const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); - - manager.announceDetectedDevice(deviceInfo); - expect(mockedEventEmitter.emit).not.toHaveBeenCalled(); - - const enabledSettings = new Settings(); - enabledSettings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); - settingsManager.getSettings.mockReturnValue(enabledSettings); + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + await connectDevice(manager, mockedEventEmitter, deviceInfo, device); await manager.onSettingsChanged(); - expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); + expect(manager.getConnectedDevices()).toHaveLength(1); }); - it('re-announces a device rejected by addDevice() only once its canonical known device gets re-enabled', async () => { + it('re-announces a device rejected by the offer only once its canonical known device gets re-enabled', async () => { // The device is detected under a preliminary id, but its final/canonical id (only // known after connecting, e.g. a serial number read during a handshake) is different. const detectionId = DeviceId.create('device-pending-2-detected'); @@ -521,8 +583,12 @@ describe('deviceManager', () => { // Simulate a provider that connected a device via the detected-device pipeline whose // final id turns out to belong to a disabled device. const device = new TestDevice(canonicalId, 'Foo', new Date(), false, new EventEmitter()); - const added = manager.addDevice(deviceInfo, device); - expect(added).toBe(false); + const result = await connectDevice(manager, mockedEventEmitter, deviceInfo, device); + expect(result.successful).toBe(false); + + // Drop the deviceDetected emit from announcing above - only the re-announce below is + // under test here, same as the original addDevice()-based version of this test. + mockClear(mockedEventEmitter); // An unrelated settings change while the canonical device is still disabled must NOT // retry it (it would if the retry were gated by the still-unknown detection id). @@ -552,7 +618,7 @@ describe('deviceManager', () => { const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); - manager.addDevice(deviceInfo, device); + await connectDevice(manager, mockedEventEmitter, deviceInfo, device); mockClear(mockedEventEmitter); diff --git a/tests/unit/device/provider/deviceProvider.spec.ts b/tests/unit/device/provider/deviceProvider.spec.ts index 66dd89c4..813a57d2 100644 --- a/tests/unit/device/provider/deviceProvider.spec.ts +++ b/tests/unit/device/provider/deviceProvider.spec.ts @@ -2,10 +2,12 @@ import { describe, expect, it, vi } from 'vitest'; import { mock } from 'vitest-mock-extended'; import EventEmitter from 'events'; import DeviceProvider from '../../../../src/device/provider/deviceProvider.js'; -import DeviceManager, { DeviceDetectionInfo } from '../../../../src/device/deviceManager.js'; +import DeviceManager, { DeviceDetectionInfo, DeviceManagerEvent } from '../../../../src/device/deviceManager.js'; import { AnyDevice } from '../../../../src/device/device.js'; import Logger from '../../../../src/logging/Logger.js'; import SettingsManager from '../../../../src/settings/settingsManager.js'; +import Settings from '../../../../src/settings/settings.js'; +import KnownDevice from '../../../../src/settings/knownDevice.js'; import { DeviceId } from '../../../../src/device/deviceId.js'; import TestDevice from '../testDevice.js'; @@ -15,15 +17,15 @@ class TestProvider extends DeviceProvider public doStopCalls = 0; public constructor(deviceManager: DeviceManager = mock()) { - super(deviceManager, new EventEmitter(), mock()); + super(deviceManager, mock()); } protected canHandleDeviceDetectionInfo(_deviceDetectionInfo: DeviceDetectionInfo): _deviceDetectionInfo is DeviceDetectionInfo { return false; } - protected createDevice(_deviceDetectionInfo: DeviceDetectionInfo): Promise { - return Promise.resolve(undefined); + protected createDevice(deviceDetectionInfo: DeviceDetectionInfo): Promise { + return Promise.resolve(new TestDevice(deviceDetectionInfo.detectionId, 'Foo', new Date(), false, new EventEmitter())); } protected override async doStart(): Promise { @@ -40,16 +42,64 @@ class TestProvider extends DeviceProvider class DetectingTestProvider extends DeviceProvider { public constructor(deviceManager: DeviceManager) { - super(deviceManager, new EventEmitter(), mock()); + super(deviceManager, mock()); } protected canHandleDeviceDetectionInfo(deviceDetectionInfo: DeviceDetectionInfo): deviceDetectionInfo is DeviceDetectionInfo { return deviceDetectionInfo.type === 'test'; } - protected createDevice(deviceDetectionInfo: DeviceDetectionInfo): Promise { + protected createDevice(deviceDetectionInfo: DeviceDetectionInfo): Promise { return Promise.resolve(new TestDevice(deviceDetectionInfo.detectionId, 'Foo', new Date(), false, new EventEmitter())); } + + // Exposes the protected getConnectedDevice() so tests can check the provider's own + // bookkeeping directly, e.g. from within a deviceManager event listener. + public hasDeviceLocally(deviceId: DeviceId): boolean { + return undefined !== this.getConnectedDevice(deviceId); + } +} + +// createDevice() resolution is controlled from outside via the injected promise, to simulate a +// slow connect attempt that's still in flight when the provider gets stopped. +class SlowCreateDeviceProvider extends DeviceProvider +{ + public constructor(deviceManager: DeviceManager, private readonly createDevicePromise: Promise) { + super(deviceManager, mock()); + } + + protected canHandleDeviceDetectionInfo(deviceDetectionInfo: DeviceDetectionInfo): deviceDetectionInfo is DeviceDetectionInfo { + return deviceDetectionInfo.type === 'test'; + } + + protected createDevice(_deviceDetectionInfo: DeviceDetectionInfo): Promise { + return this.createDevicePromise; + } +} + +// Tracks onConnectFailed() calls so tests can assert it does/doesn't run for a given rejection. +class TrackingTestProvider extends DeviceProvider +{ + public onConnectFailedCalls = 0; + + public constructor( + deviceManager: DeviceManager, + private readonly createDeviceFn: (deviceDetectionInfo: DeviceDetectionInfo) => Promise + ) { + super(deviceManager, mock()); + } + + protected canHandleDeviceDetectionInfo(deviceDetectionInfo: DeviceDetectionInfo): deviceDetectionInfo is DeviceDetectionInfo { + return deviceDetectionInfo.type === 'test'; + } + + protected createDevice(deviceDetectionInfo: DeviceDetectionInfo): Promise { + return this.createDeviceFn(deviceDetectionInfo); + } + + protected override async onConnectFailed(_deviceDetectionInfo: DeviceDetectionInfo): Promise { + this.onConnectFailedCalls++; + } } describe('DeviceProvider', () => { @@ -125,4 +175,113 @@ describe('DeviceProvider', () => { await vi.waitFor(() => expect(deviceManager.getConnectedDevices()).toHaveLength(1)); }); }); + + describe('handleDeviceDetection', () => { + it('has the device in its own connected list by the time deviceManager emits deviceConnected', async () => { + const settingsManager = mock(); + settingsManager.getSettings.mockReturnValue(undefined); + + const logger = mock(); + logger.child.mockReturnValue(logger); + + const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + const provider = new DetectingTestProvider(deviceManager); + await provider.start(); + + const deviceId = DeviceId.create('device-race'); + let sawItLocallyOnConnect = false; + + deviceManager.on(DeviceManagerEvent.deviceConnected, (device) => { + if (device.getDeviceId === deviceId) { + sawItLocallyOnConnect = provider.hasDeviceLocally(deviceId); + } + }); + + deviceManager.announceDetectedDevice({ type: 'test', detectionId: deviceId }); + + await vi.waitFor(() => expect(deviceManager.getConnectedDevices()).toHaveLength(1)); + + expect(sawItLocallyOnConnect).toBe(true); + }); + + it('closes and does not register a device that connects after the provider was stopped', async () => { + const settingsManager = mock(); + settingsManager.getSettings.mockReturnValue(undefined); + + const logger = mock(); + logger.child.mockReturnValue(logger); + + const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + + let resolveCreateDevice!: (device: AnyDevice) => void; + const createDevicePromise = new Promise((resolve) => { resolveCreateDevice = resolve; }); + + const provider = new SlowCreateDeviceProvider(deviceManager, createDevicePromise); + await provider.start(); + + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DeviceId.create('device-stopped') }); + + // Provider is stopped while createDevice() is still pending + await provider.stop(); + + const device = new TestDevice(DeviceId.create('device-stopped'), 'Foo', new Date(), false, new EventEmitter()); + const closeSpy = vi.spyOn(device, 'close'); + + resolveCreateDevice(device); + + await vi.waitFor(() => expect(closeSpy).toHaveBeenCalled()); + + expect(deviceManager.getConnectedDevices()).toHaveLength(0); + }); + + it('calls onConnectFailed when the offer itself throws', async () => { + const settingsManager = mock(); + settingsManager.getSettings.mockReturnValue(undefined); + + const logger = mock(); + logger.child.mockReturnValue(logger); + + const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + const provider = new TrackingTestProvider(deviceManager, () => Promise.reject(new Error('connect failed'))); + + await provider.start(); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DeviceId.create('device-throw') }); + + await vi.waitFor(() => expect(provider.onConnectFailedCalls).toBe(1)); + }); + + it('does not call onConnectFailed when the device is rejected for being disabled', async () => { + const deviceId = DeviceId.create('device-disabled-oncf'); + const settings = new Settings(); + settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, false)); + + const settingsManager = mock(); + settingsManager.getSettings.mockReturnValue(settings); + + const logger = mock(); + logger.child.mockReturnValue(logger); + + const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + + let closeSpy: ReturnType | undefined; + const provider = new TrackingTestProvider( + deviceManager, + (deviceDetectionInfo) => { + const device = new TestDevice(deviceDetectionInfo.detectionId, 'Foo', new Date(), false, new EventEmitter()); + closeSpy = vi.spyOn(device, 'close'); + return Promise.resolve(device); + } + ); + + await provider.start(); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: deviceId }); + + // Disabled devices are closed internally by offerDevice() once rejected - wait for + // that deterministically instead of a fixed sleep. + await vi.waitFor(() => expect(closeSpy).toHaveBeenCalled()); + + expect(deviceManager.getConnectedDevices()).toHaveLength(0); + expect(provider.onConnectFailedCalls).toBe(0); + }); + }); }); diff --git a/tests/unit/device/provider/deviceProviderManager.spec.ts b/tests/unit/device/provider/deviceProviderManager.spec.ts index c1e6a99d..f68babb7 100644 --- a/tests/unit/device/provider/deviceProviderManager.spec.ts +++ b/tests/unit/device/provider/deviceProviderManager.spec.ts @@ -22,7 +22,7 @@ class RecordingDeviceProvider extends DeviceProvider = Promise.resolve(); public constructor() { - super(mock(), new EventEmitter(), mock()); + super(mock(), mock()); } public setStartGate(gate: Promise): void { @@ -49,8 +49,8 @@ class RecordingDeviceProvider extends DeviceProvider { - return Promise.resolve(undefined); + protected createDevice(_deviceDetectionInfo: DeviceDetectionInfo): Promise { + return Promise.resolve(mock()); } }