diff --git a/docker/dev/node/Dockerfile b/docker/dev/node/Dockerfile index 9be2978b..cb7cdaf6 100644 --- a/docker/dev/node/Dockerfile +++ b/docker/dev/node/Dockerfile @@ -1,4 +1,4 @@ -FROM node:18-alpine +FROM node:24-alpine RUN apk --no-cache upgrade && \ apk --no-cache add bash git sudo openssh make diff --git a/eslint.config.ts b/eslint.config.ts index 46789984..64d32de9 100644 --- a/eslint.config.ts +++ b/eslint.config.ts @@ -90,7 +90,12 @@ export default [ "@typescript-eslint/no-empty-function": "error", "@typescript-eslint/no-empty-interface": "off", "@typescript-eslint/no-explicit-any": "off", - "@typescript-eslint/no-floating-promises": "warn", + "@typescript-eslint/no-floating-promises": [ + "warn", + { + "checkThenables": true + } + ], "@typescript-eslint/unbound-method": "error", "@typescript-eslint/no-misused-promises": "error", "@typescript-eslint/no-misused-new": "error", diff --git a/package-lock.json b/package-lock.json index 2ef25135..d0d28aeb 100644 --- a/package-lock.json +++ b/package-lock.json @@ -9,10 +9,11 @@ "@sinclair/typebox": "^0.34.47", "@stoprocent/noble": "^2.3.17", "@timesplinter/pimple": "^2.1.1", - "@timesplinter/sequential-task-queue": "^1.3.1", + "@timesplinter/sequential-task-queue": "^1.4.0", "ajv": "^8.17.1", "ajv-formats": "^3.0.1", "buttplug": "^3.2.2", + "chokidar": "^5.0.0", "class-transformer": "^0.5.1", "cors": "^2.8.5", "dotenv": "^17.2.0", @@ -3172,9 +3173,9 @@ "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==", + "version": "1.4.0", + "resolved": "https://registry.npmjs.org/@timesplinter/sequential-task-queue/-/sequential-task-queue-1.4.0.tgz", + "integrity": "sha512-yXVO3Va7d2UiWyw7ucrsFAbTOtDIrV3DghmSMCg0jRkJNpMW3v65whddK3IXEuaJxpMR2E7UQ7AXLx0410Gdtw==", "license": "MIT" }, "node_modules/@tybys/wasm-util": { @@ -3848,6 +3849,8 @@ }, "node_modules/anymatch": { "version": "3.1.3", + "resolved": "https://registry.npmjs.org/anymatch/-/anymatch-3.1.3.tgz", + "integrity": "sha512-KMReFUr0B4t+D+OBkjR3KYqvocp2XaSzO55UcB6mgQMd3KbcE+mWTyvVV7D/zsdEbNnV6acZUutkiHQXvTr1Rw==", "dev": true, "license": "ISC", "dependencies": { @@ -4075,6 +4078,8 @@ }, "node_modules/binary-extensions": { "version": "2.3.0", + "resolved": "https://registry.npmjs.org/binary-extensions/-/binary-extensions-2.3.0.tgz", + "integrity": "sha512-Ceh+7ox5qe7LJuLHoY0feh3pHuUDHAcRUeyL2VYghZwfpkNIy/+8Ocg0a3UuSoYzavmylwuLWQOf3hl0jjMMIw==", "dev": true, "license": "MIT", "engines": { @@ -4313,37 +4318,18 @@ } }, "node_modules/chokidar": { - "version": "3.6.0", - "dev": true, + "version": "5.0.0", + "resolved": "https://registry.npmjs.org/chokidar/-/chokidar-5.0.0.tgz", + "integrity": "sha512-TQMmc3w+5AxjpL8iIiwebF73dRDF4fBIieAqGn9RGCWaEVwQ6Fb2cGe31Yns0RRIzii5goJ1Y7xbMwo1TxMplw==", "license": "MIT", "dependencies": { - "anymatch": "~3.1.2", - "braces": "~3.0.2", - "glob-parent": "~5.1.2", - "is-binary-path": "~2.1.0", - "is-glob": "~4.0.1", - "normalize-path": "~3.0.0", - "readdirp": "~3.6.0" + "readdirp": "^5.0.0" }, "engines": { - "node": ">= 8.10.0" + "node": ">= 20.19.0" }, "funding": { "url": "https://paulmillr.com/funding/" - }, - "optionalDependencies": { - "fsevents": "~2.3.2" - } - }, - "node_modules/chokidar/node_modules/glob-parent": { - "version": "5.1.2", - "dev": true, - "license": "ISC", - "dependencies": { - "is-glob": "^4.0.1" - }, - "engines": { - "node": ">= 6" } }, "node_modules/ci-info": { @@ -6005,6 +5991,8 @@ }, "node_modules/is-binary-path": { "version": "2.1.0", + "resolved": "https://registry.npmjs.org/is-binary-path/-/is-binary-path-2.1.0.tgz", + "integrity": "sha512-ZMERYes6pDydyuGidse7OsHxtbI7WVeUEozgR/g7rd0xUimYNlvZRE/K2MgZTjWy725IfelLeVcEM97mmtRGXw==", "dev": true, "license": "MIT", "dependencies": { @@ -7238,6 +7226,44 @@ "url": "https://opencollective.com/nodemon" } }, + "node_modules/nodemon/node_modules/chokidar": { + "version": "3.6.0", + "resolved": "https://registry.npmjs.org/chokidar/-/chokidar-3.6.0.tgz", + "integrity": "sha512-7VT13fmjotKpGipCW9JEQAusEPE+Ei8nl6/g4FBAmIm0GOOLMua9NDDo/DWp0ZAxCr3cPq5ZpBqmPAQgDda2Pw==", + "dev": true, + "license": "MIT", + "dependencies": { + "anymatch": "~3.1.2", + "braces": "~3.0.2", + "glob-parent": "~5.1.2", + "is-binary-path": "~2.1.0", + "is-glob": "~4.0.1", + "normalize-path": "~3.0.0", + "readdirp": "~3.6.0" + }, + "engines": { + "node": ">= 8.10.0" + }, + "funding": { + "url": "https://paulmillr.com/funding/" + }, + "optionalDependencies": { + "fsevents": "~2.3.2" + } + }, + "node_modules/nodemon/node_modules/glob-parent": { + "version": "5.1.2", + "resolved": "https://registry.npmjs.org/glob-parent/-/glob-parent-5.1.2.tgz", + "integrity": "sha512-AOIgSQCepiJYwP3ARnGx+5VnTu2HBYdzbGP45eLw1vr3zB3vZLeyed1sC9hnbcOc9/SrMyM5RPQrkGz4aS9Zow==", + "dev": true, + "license": "ISC", + "dependencies": { + "is-glob": "^4.0.1" + }, + "engines": { + "node": ">= 6" + } + }, "node_modules/nodemon/node_modules/has-flag": { "version": "3.0.0", "dev": true, @@ -7246,6 +7272,19 @@ "node": ">=4" } }, + "node_modules/nodemon/node_modules/readdirp": { + "version": "3.6.0", + "resolved": "https://registry.npmjs.org/readdirp/-/readdirp-3.6.0.tgz", + "integrity": "sha512-hOS089on8RduqdbhvQ5Z37A0ESjsqz6qnRcffsMU3495FuTdqSm+7bhJ29JvIOsBDEEnan5DPu9t3To9VRlMzA==", + "dev": true, + "license": "MIT", + "dependencies": { + "picomatch": "^2.2.1" + }, + "engines": { + "node": ">=8.10.0" + } + }, "node_modules/nodemon/node_modules/semver": { "version": "7.8.5", "resolved": "https://registry.npmjs.org/semver/-/semver-7.8.5.tgz", @@ -7283,6 +7322,8 @@ }, "node_modules/normalize-path": { "version": "3.0.0", + "resolved": "https://registry.npmjs.org/normalize-path/-/normalize-path-3.0.0.tgz", + "integrity": "sha512-6eZs5Ls3WtCisHWp9S2GUy8dqkpGi4BVSz3GaqiE6ezub0512ESztXUwUB6C6IKbQkY2Pnb/mD4WYojCRwcwLA==", "dev": true, "license": "MIT", "engines": { @@ -7914,14 +7955,16 @@ } }, "node_modules/readdirp": { - "version": "3.6.0", - "dev": true, + "version": "5.0.0", + "resolved": "https://registry.npmjs.org/readdirp/-/readdirp-5.0.0.tgz", + "integrity": "sha512-9u/XQ1pvrQtYyMpZe7DXKv2p5CNvyVwzUB6uhLAnQwHMSgKMBR62lc7AHljaeteeHXn11XTAaLLUVZYVZyuRBQ==", "license": "MIT", - "dependencies": { - "picomatch": "^2.2.1" - }, "engines": { - "node": ">=8.10.0" + "node": ">= 20.19.0" + }, + "funding": { + "type": "individual", + "url": "https://paulmillr.com/funding/" } }, "node_modules/real-require": { diff --git a/package.json b/package.json index 1bca2260..dcf7bd8f 100644 --- a/package.json +++ b/package.json @@ -10,6 +10,7 @@ "ajv": "^8.17.1", "ajv-formats": "^3.0.1", "buttplug": "^3.2.2", + "chokidar": "^5.0.0", "class-transformer": "^0.5.1", "cors": "^2.8.5", "dotenv": "^17.2.0", @@ -22,7 +23,7 @@ "read-last-lines": "^1.8.0", "reflect-metadata": "^0.1.13", "say": "^0.16.0", - "@timesplinter/sequential-task-queue": "^1.3.1", + "@timesplinter/sequential-task-queue": "^1.4.0", "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/app.ts b/src/app.ts index 811b1633..809509a2 100644 --- a/src/app.ts +++ b/src/app.ts @@ -207,6 +207,10 @@ export const createApp = (container: Container, options: AppOptions) .use(express.text()) ; + const settingsManager = container.get('settings.manager'); + settingsManager.load(); + settingsManager.startWatching(); + configureRoutes(app, container); configureWebsocket(websocketServer, container); startDeviceProviders(container); @@ -262,14 +266,24 @@ export const createApp = (container: Container, options: AppOptions) const logger = container.get('logger.default'); logger.info('Shutting down...'); - await container.get('automation.scriptRuntime').stop(); + try { + await container.get('settings.manager').stopWatching(); + } catch (e: unknown) { + logError(logger, 'Failed to stop settings file watcher during shutdown', e); + } try { - await container.get('device.provider.manager').stopProviders(); + await container.get('automation.scriptRuntime').stop(); } catch (e: unknown) { - logError(logger, 'Failed to stop device providers during shutdown', e); + logError(logger, 'Failed to stop automation script runtime during shutdown', e); } + try { + await container.get('device.provider.manager').stopProviders(); + } catch (e: unknown) { + logError(logger, 'Failed to stop device providers during shutdown', e); + } + container.get('health.metricsCollector').stop(); await websocketServer.close(); diff --git a/src/device/bleDevice.ts b/src/device/bleDevice.ts index 4ffb9a1b..84ebfc18 100644 --- a/src/device/bleDevice.ts +++ b/src/device/bleDevice.ts @@ -26,8 +26,6 @@ export default abstract class BleDevice< @Expose() private rssi: number; - protected logger: Logger; - protected constructor( deviceId: DeviceId, deviceName: string, @@ -40,9 +38,7 @@ export default abstract class BleDevice< eventEmitter: EventEmitter, logger: Logger, ) { - super(deviceId, deviceName, provider, connectedSince, controllable, attributes, config, eventEmitter); - - this.logger = logger.child({ name: this.constructor.name }); + super(deviceId, deviceName, provider, connectedSince, controllable, attributes, config, eventEmitter, logger); this.peripheral = peripheral; this.rssi = peripheral.rssi; diff --git a/src/device/detectedDeviceOfferQueue.ts b/src/device/detectedDeviceOfferQueue.ts index 681ffde2..9f73eb30 100644 --- a/src/device/detectedDeviceOfferQueue.ts +++ b/src/device/detectedDeviceOfferQueue.ts @@ -4,6 +4,7 @@ import { DeviceDetectionInfo } from './deviceManager.js'; import DeviceOfferRejectedError from './deviceOfferRejectedError.js'; import Logger from '../logging/Logger.js'; import { logError } from '../util/error.js'; +import { DetectionId } from './deviceId.js'; export type OfferResult = | { successful: true, device: D } @@ -13,7 +14,7 @@ type DeviceOffer = (cancellationToken: CancellationToken) = export default class DetectedDeviceOfferQueue { - private readonly queues: Map = new Map(); + private readonly queues: Map = new Map(); private readonly logger: Logger; @@ -21,7 +22,7 @@ export default class DetectedDeviceOfferQueue this.logger = logger; } - private getOrCreateQueue(detectionId: string): SequentialTaskQueue + private getOrCreateQueue(detectionId: DetectionId): SequentialTaskQueue { let queue = this.queues.get(detectionId); @@ -58,7 +59,7 @@ export default class DetectedDeviceOfferQueue const task = queue.push((cancellationToken: CancellationToken) => this.runOffer(deviceOffer, cancellationToken)); - return Promise.resolve(task.then( + return task.then( (result: OfferResult): OfferResult => { if (result.successful) { // Reject every other still-queued offer for this detection id without them @@ -75,7 +76,7 @@ export default class DetectedDeviceOfferQueue successful: false, reason: reason, }) - )); + ); } private async runOffer( @@ -111,12 +112,12 @@ export default class DetectedDeviceOfferQueue * callers that need "is a fresh announce still blocked by a past revoke" must call * dropIfRevoked() first. */ - public has(detectionId: string): boolean + public has(detectionId: DetectionId): boolean { return this.queues.has(detectionId); } - public dropIfRevoked(detectionId: string): void + public dropIfRevoked(detectionId: DetectionId): void { const queue = this.queues.get(detectionId); @@ -125,7 +126,7 @@ export default class DetectedDeviceOfferQueue } } - private close(detectionId: string, reason: DeviceOfferRejectedError): void + private close(detectionId: DetectionId, reason: DeviceOfferRejectedError): void { const queue = this.queues.get(detectionId); @@ -136,7 +137,7 @@ export default class DetectedDeviceOfferQueue this.queues.delete(detectionId); } - public revoke(detectionId: string, reason: DeviceOfferRejectedError): void + public revoke(detectionId: DetectionId, reason: DeviceOfferRejectedError): void { const queue = this.getOrCreateQueue(detectionId); diff --git a/src/device/device.ts b/src/device/device.ts index 8df695dd..4259df04 100644 --- a/src/device/device.ts +++ b/src/device/device.ts @@ -6,6 +6,7 @@ import { EventEmitter } from 'events'; import type { DeviceId } from './deviceId.js'; import type { JsonObject } from '../types.js'; import { DropFirst } from '../types.js'; +import Logger from '../logging/Logger.js'; // An attribute value can be DeviceAttribute or undefined because we want to allow Partial<> export type DeviceAttributes = Record; @@ -96,6 +97,10 @@ export default abstract class Device< private eventEmitter: EventEmitter; + private closePromise?: Promise; + + protected readonly logger: Logger; + protected constructor( deviceId: DeviceId, deviceName: string, @@ -104,7 +109,8 @@ export default abstract class Device< controllable: boolean, attributes: TAttributes, config: TConfig, - eventEmitter: EventEmitter + eventEmitter: EventEmitter, + logger: Logger ) { this.deviceId = deviceId; this.deviceName = deviceName; @@ -114,6 +120,7 @@ export default abstract class Device< this.attributes = attributes; this.config = config; this.eventEmitter = eventEmitter; + this.logger = logger.child({ name: `${new.target.name}.${deviceId}` }); this.state = DeviceState.ready; } @@ -143,8 +150,8 @@ export default abstract class Device< public async refresh(): Promise { - if (this.state === DeviceState.closed) { - throw new Error('Cannot refresh device as it is closed'); + if (this.state === DeviceState.closed || this.state === DeviceState.closing) { + throw new Error('Cannot refresh device as it is closed or closing'); } await this.doRefresh(); @@ -178,10 +185,18 @@ export default abstract class Device< public async close(): Promise { - if (this.state === DeviceState.closed) { - return + if (undefined !== this.closePromise) { + return this.closePromise; } + this.state = DeviceState.closing; + this.closePromise = this.performClose(); + + return this.closePromise; + } + + private async performClose(): Promise + { try { await this.doClose(); } finally { diff --git a/src/device/deviceId.ts b/src/device/deviceId.ts index de4ff5a0..d55d0802 100644 --- a/src/device/deviceId.ts +++ b/src/device/deviceId.ts @@ -2,12 +2,29 @@ import { v5 as uuidv5 } from 'uuid'; const DEVICE_NAMESPACE = '1e0758c9-799d-40b5-b2fc-63f1e66afb76'; const deviceIdSymbol = Symbol(); +const detectionIdSymbol = Symbol(); export type DeviceId = string & { [deviceIdSymbol]: never } +export type DetectionId = string & { [detectionIdSymbol]: never } export const DeviceId = { create: (seed: string): DeviceId => { // eslint-disable-next-line @typescript-eslint/consistent-type-assertions return uuidv5(seed, DEVICE_NAMESPACE).toString() as DeviceId; - } + }, + fromDetectionId: (detectionId: DetectionId): DeviceId => { + // eslint-disable-next-line @typescript-eslint/consistent-type-assertions + return detectionId as unknown as DeviceId; + }, +} + +export const DetectionId = { + create: (seed: string): DetectionId => { + // eslint-disable-next-line @typescript-eslint/consistent-type-assertions + return uuidv5(seed, DEVICE_NAMESPACE).toString() as DetectionId; + }, + fromDeviceId: (deviceId: DeviceId): DetectionId => { + // eslint-disable-next-line @typescript-eslint/consistent-type-assertions + return deviceId as unknown as DetectionId; + }, } diff --git a/src/device/deviceManager.ts b/src/device/deviceManager.ts index c625cbaa..ffeab52c 100644 --- a/src/device/deviceManager.ts +++ b/src/device/deviceManager.ts @@ -5,14 +5,14 @@ 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 { DeviceId, DetectionId } 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; - detectionId: DeviceId; + detectionId: DetectionId; }; export enum DeviceManagerEvent { @@ -37,6 +37,8 @@ type DeviceManagerEventMap = { [DeviceManagerEvent.deviceNotification]: [device: AnyDevice, notification: DeviceNotification]; } +type ConnectedDevice = { device: AnyDevice, deviceDetectionInfo: DeviceDetectionInfo }; + export default class DeviceManager { private readonly eventEmitter: EventEmitter; @@ -45,24 +47,22 @@ export default class DeviceManager private readonly offerQueue: DetectedDeviceOfferQueue; - private readonly connectedDevices: Map; - private readonly settingsManager: SettingsManager; - private readonly detectedDisabledDevices: Map = new Map(); + private readonly detectedDisabledDevices: Map = new Map(); + + private readonly connectedDevices: Map = new Map(); // Serializes onSettingsChanged() runs so rapid settings changes don't interleave private readonly settingsChangeQueue: SequentialTaskQueue = new SequentialTaskQueue(); public constructor( eventEmitter: EventEmitter, - connectedDevices: Map, settingsManager: SettingsManager, logger: Logger ) { this.eventEmitter = eventEmitter; this.logger = logger.child({ name: DeviceManager.name }); - this.connectedDevices = connectedDevices; this.settingsManager = settingsManager; this.offerQueue = new DetectedDeviceOfferQueue(this.logger); } @@ -71,6 +71,17 @@ export default class DeviceManager return this.settingsManager.getSettings()?.getKnownDeviceById(deviceId)?.enabled ?? true; } + private isDetectedDeviceAlreadyConnected(detectionId: DetectionId): boolean + { + for (const { deviceDetectionInfo } of this.connectedDevices.values()) { + if (deviceDetectionInfo.detectionId === detectionId) { + return true; + } + } + + return false; + } + public announceDetectedDevice(deviceDetectionInfo: DeviceDetectionInfo): void { this.offerQueue.dropIfRevoked(deviceDetectionInfo.detectionId); @@ -79,7 +90,7 @@ export default class DeviceManager return; } - if (this.connectedDevices.has(deviceDetectionInfo.detectionId)) { + if (this.isDetectedDeviceAlreadyConnected(deviceDetectionInfo.detectionId)) { this.logger.debug(`Device with id '${deviceDetectionInfo.detectionId}' is already connected, not announcing it as detected`); return; } @@ -107,7 +118,7 @@ export default class DeviceManager public async offerDevice(deviceDetectionInfo: DeviceDetectionInfo, deviceOffer: () => Promise): Promise> { - if (this.connectedDevices.has(deviceDetectionInfo.detectionId)) { + if (this.isDetectedDeviceAlreadyConnected(deviceDetectionInfo.detectionId)) { return { successful: false, reason: new DeviceOfferRejectedError('Device is already connected') }; } @@ -140,13 +151,13 @@ export default class DeviceManager }); if (result.successful) { - this.registerDevice(result.device); + this.registerDevice(result.device, deviceDetectionInfo); } return result; } - private registerDevice(device: AnyDevice): void + private registerDevice(device: AnyDevice, deviceDetectionInfo: DeviceDetectionInfo): void { device.on(DeviceEvent.deviceRefreshed, (d) => this.eventEmitter.emit(DeviceManagerEvent.deviceRefreshed, d)); device.on(DeviceEvent.deviceDisconnected, (d) => { @@ -157,7 +168,7 @@ export default class DeviceManager this.initDeviceRefresher(device); - this.connectedDevices.set(device.getDeviceId, device); + this.connectedDevices.set(device.getDeviceId, { device, deviceDetectionInfo }); this.eventEmitter.emit(DeviceManagerEvent.deviceConnected, device); } @@ -167,18 +178,23 @@ export default class DeviceManager } private async applySettingsChange(): Promise { - for (const device of this.connectedDevices.values()) { + for (const { device, deviceDetectionInfo } of this.connectedDevices.values()) { if (this.isDeviceEnabled(device.getDeviceId)) { continue; } this.logger.info(`Closing device '${device.getDeviceId}' since it has been disabled`); - try { - await device.close(); - } catch (e: unknown) { - logError(this.logger, `Failed to close device '${device.getDeviceId}'`, e); - } + const deviceReleased = device.close() + .catch((e: unknown) => logError(this.logger, `Failed to close device '${device.getDeviceId}'`, e)); + + // Registered before awaiting deviceReleased below, so re-enabling this known device is + // picked up even if its provider never stops/restarts throughout (e.g. only this one + // known device was disabled, not its whole device source) - without it, nothing else + // would ever re-announce it. + this.detectedDisabledDevices.set(deviceDetectionInfo.detectionId, { deviceDetectionInfo, canonicalId: device.getDeviceId, deviceReleased }); + + await deviceReleased; } for (const [detectionId, disabledDetectedDevice] of [...this.detectedDisabledDevices]) { @@ -197,14 +213,12 @@ export default class DeviceManager public getConnectedDevices(): AnyDevice[] { - return Array.from(this.connectedDevices.values()); + return Array.from(this.connectedDevices.values(), (entry) => entry.device); } - public getConnectedDevice(deviceId: string): AnyDevice|null + public getConnectedDevice(deviceId: DeviceId): AnyDevice|null { - const device = this.connectedDevices.get(deviceId); - - return undefined !== device ? device : null; + return this.connectedDevices.get(deviceId)?.device ?? null; } public on( @@ -229,7 +243,7 @@ export default class DeviceManager let closeError: unknown; - for (const [, device] of this.connectedDevices) { + for (const { device } of this.connectedDevices.values()) { try { await device.close(); } catch (e: unknown) { @@ -256,8 +270,8 @@ export default class DeviceManager } const deviceRefresher = async (): Promise => { - if (device.getState === DeviceState.busy) { - this.logger.trace(`Device not refreshed since it's currently busy: ${device.getDeviceId}`); + if (device.getState === DeviceState.busy || device.getState === DeviceState.closing || device.getState === DeviceState.closed) { + this.logger.trace(`Device not refreshed since it's currently ${device.getState}: ${device.getDeviceId}`); return; } diff --git a/src/device/deviceState.ts b/src/device/deviceState.ts index 782de18a..4525fe38 100644 --- a/src/device/deviceState.ts +++ b/src/device/deviceState.ts @@ -1,6 +1,7 @@ enum DeviceState { ready = 'READY', busy = 'BUSY', + closing = 'CLOSING', error = 'ERROR', closed = 'CLOSED', } diff --git a/src/device/peripheralDevice.ts b/src/device/peripheralDevice.ts index 28d1d063..dcc5eac3 100644 --- a/src/device/peripheralDevice.ts +++ b/src/device/peripheralDevice.ts @@ -4,6 +4,8 @@ import DeviceProtocol, { MessageWithResponse } from './protocol/deviceProtocol.j import { AnyDeviceConfig, NoDeviceConfig } from './deviceConfig.js'; import EventEmitter from 'events'; import { DeviceId } from './deviceId.js'; +import Logger from '../logging/Logger.js'; +import { logError } from '../util/error.js'; export type AnyPeripheralDevice = WithUntypedAttributes>>>; @@ -28,14 +30,17 @@ export default abstract class PeripheralDevice< transport: BidirectionalDeviceTransport, attributes: TAttributes, config: TConfig, - eventEmitter: EventEmitter + eventEmitter: EventEmitter, + logger: Logger ) { - super(deviceId, deviceName, provider, connectedSince, controllable, attributes, config, eventEmitter); + super(deviceId, deviceName, provider, connectedSince, controllable, attributes, config, eventEmitter, logger); this.protocol = protocol; this.transport = transport; - this.transport.onClose(async () => await this.close()); + this.transport.onClose(async () => { + this.close().catch((err: unknown) => logError(this.logger, 'Error closing device after transport close', err)); + }); } public getTransport(): BidirectionalDeviceTransport { diff --git a/src/device/protocol/airotic/airoticDeviceFactory.ts b/src/device/protocol/airotic/airoticDeviceFactory.ts index e8a7627b..a3722570 100644 --- a/src/device/protocol/airotic/airoticDeviceFactory.ts +++ b/src/device/protocol/airotic/airoticDeviceFactory.ts @@ -9,7 +9,7 @@ import BoolDeviceAttribute from '../../attribute/boolDeviceAttribute.js'; import FloatDeviceAttribute from '../../attribute/floatDeviceAttribute.js'; import BleUartDeviceTransport from '../../transport/bleDeviceTransport.js'; import { hsvByteToRgb } from '../../../util/color.js'; -import { DeviceId } from '../../deviceId.js'; +import { DeviceId, DetectionId } from '../../deviceId.js'; import AiroticDevice, { AiroticDeviceAttributes } from './airoticDevice.js'; import AiroticProtocol from './airoticProtocol.js'; import MessageResponseHandler from '../messageResponseHandler.js'; @@ -37,12 +37,14 @@ export default class AiroticDeviceFactory } public create( - deviceId: DeviceId, + detectionId: DetectionId, peripheral: Peripheral, transport: BleUartDeviceTransport, messageResponseHandler: MessageResponseHandler, provider: string ): AiroticDevice { + const deviceId = DeviceId.fromDetectionId(detectionId); + const knownDevice = this.knownDeviceRegistry.resolve( deviceId, 'airotic', diff --git a/src/device/protocol/buttplugIo/buttplugIoDevice.ts b/src/device/protocol/buttplugIo/buttplugIoDevice.ts index b087cbba..cf96e2ce 100644 --- a/src/device/protocol/buttplugIo/buttplugIoDevice.ts +++ b/src/device/protocol/buttplugIo/buttplugIoDevice.ts @@ -44,14 +44,13 @@ export default class ButtplugIoDevice extends Device eventEmitter: EventEmitter, logger: Logger ) { - super(deviceId, deviceName, provider, connectedSince, true, attributes, {}, eventEmitter); + super(deviceId, deviceName, provider, connectedSince, true, attributes, {}, eventEmitter, logger); this.buttplugClientDevice = buttplugClientDevice; this.deviceModel = deviceModel; - const deviceLogger = logger.child({ name: ButtplugIoDevice.name }); this.deviceRemovedHandler = asyncHandler( async () => { await this.close(); }, - (e: unknown) => logError(deviceLogger, `Failed to close removed device '${deviceId}'`, e) + (e: unknown) => logError(this.logger, `Failed to close removed device '${deviceId}'`, e) ); this.buttplugClientDevice.on('deviceremoved', this.deviceRemovedHandler); } diff --git a/src/device/protocol/buttplugIo/buttplugIoDeviceFactory.ts b/src/device/protocol/buttplugIo/buttplugIoDeviceFactory.ts index d51c3ffb..fa8462a4 100644 --- a/src/device/protocol/buttplugIo/buttplugIoDeviceFactory.ts +++ b/src/device/protocol/buttplugIo/buttplugIoDeviceFactory.ts @@ -10,7 +10,7 @@ import DateFactory from '../../../factory/dateFactory.js'; import { Int } from '../../../util/numbers.js'; import IntDeviceAttribute from '../../attribute/intDeviceAttribute.js'; import EventEmitterFactory from '../../../factory/eventEmitterFactory.js'; -import { DeviceId } from '../../deviceId.js'; +import { DeviceId, DetectionId } from '../../deviceId.js'; export default class ButtplugIoDeviceFactory @@ -36,8 +36,8 @@ export default class ButtplugIoDeviceFactory this.logger = logger; } - public create(deviceId: DeviceId, buttplugDevice: ButtplugClientDevice, provider: string): ButtplugIoDevice { - const knownDevice = this.resolveKnownDevice(deviceId, buttplugDevice, provider); + public create(detectionId: DetectionId, buttplugDevice: ButtplugClientDevice, provider: string): ButtplugIoDevice { + const knownDevice = this.resolveKnownDevice(DeviceId.fromDetectionId(detectionId), buttplugDevice, provider); const deviceAttrs = ButtplugIoDeviceFactory.parseDeviceAttributes(buttplugDevice); diff --git a/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProvider.ts b/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProvider.ts index c9670854..f13357d4 100644 --- a/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProvider.ts +++ b/src/device/protocol/buttplugIo/buttplugIoWebsocketDeviceProvider.ts @@ -8,7 +8,7 @@ import SlvCtrlPlusButtplugWebsocketClientConnector from './slvCtrlPlusButtplugWe import DeviceManager, { DeviceDetectionInfo } from '../../deviceManager.js'; import { logError } from '../../../util/error.js'; import { hasProperty } from '../../../util/objects.js'; -import { DeviceId } from '../../deviceId.js'; +import { DeviceId, DetectionId } from '../../deviceId.js'; export type ButtplugIoDeviceDetectionInfo = DeviceDetectionInfo & { type: 'buttplugIo'; @@ -170,7 +170,7 @@ export default class ButtplugIoWebsocketDeviceProvider extends DeviceProvider< const nameString = buttplugDevice.name.replace(/[^a-zA-Z0-9]/g, ''); const deviceId = DeviceId.create(this.useDeviceNameAsId ? `buttplugio-${nameString}` : `buttplugio-${buttplugDevice.index}`); - return { type: 'buttplugIo', detectionId: deviceId, buttplugClientDevice: buttplugDevice }; + return { type: 'buttplugIo', detectionId: DetectionId.fromDeviceId(deviceId), buttplugClientDevice: buttplugDevice }; } private announceButtplugIoDevice(buttplugDevice: ButtplugClientDevice): void { diff --git a/src/device/protocol/estim2b/estim2bDevice.ts b/src/device/protocol/estim2b/estim2bDevice.ts index b28e89f5..e089085f 100644 --- a/src/device/protocol/estim2b/estim2bDevice.ts +++ b/src/device/protocol/estim2b/estim2bDevice.ts @@ -36,8 +36,6 @@ export default class EStim2bDevice extends PeripheralDevice { const attributes = this.getAttributes(initialStatus); - const knownDevice = this.knownDeviceRegistry.resolve(deviceId, 'estim2b', provider); - - // KnownDevice is not persisted as we cannot determine a unique device id for the estim2b device, - // so we cannot reliably identify it on future connections. + const knownDevice = this.knownDeviceRegistry.resolve(DeviceId.fromDetectionId(detectionId), 'estim2b', provider); return new Estim2bDevice( knownDevice.id, diff --git a/src/device/protocol/slvCtrlPlus/slvCtrlPlusDevice.ts b/src/device/protocol/slvCtrlPlus/slvCtrlPlusDevice.ts index c4069047..0b01aa38 100644 --- a/src/device/protocol/slvCtrlPlus/slvCtrlPlusDevice.ts +++ b/src/device/protocol/slvCtrlPlus/slvCtrlPlusDevice.ts @@ -19,8 +19,6 @@ export default abstract class SlvCtrlPlusDevice< TNotifications extends DeviceNotifications = NoDeviceNotifications, TConfig extends AnyDeviceConfig = NoDeviceConfig, > extends PeripheralDevice { - protected readonly logger: Logger; - protected constructor( deviceId: DeviceId, deviceName: string, @@ -34,9 +32,7 @@ export default abstract class SlvCtrlPlusDevice< eventEmitter: EventEmitter, logger: Logger, ) { - super(deviceId, deviceName, provider, connectedSince, controllable, protocol, transport, attributes, config, eventEmitter); - - this.logger = logger; + super(deviceId, deviceName, provider, connectedSince, controllable, protocol, transport, attributes, config, eventEmitter, logger); } protected async send(command: SlvCtrlProtocolCommand): Promise diff --git a/src/device/protocol/slvCtrlPlus/slvCtrlPlusDeviceFactory.ts b/src/device/protocol/slvCtrlPlus/slvCtrlPlusDeviceFactory.ts index a13de621..3dfe62fd 100644 --- a/src/device/protocol/slvCtrlPlus/slvCtrlPlusDeviceFactory.ts +++ b/src/device/protocol/slvCtrlPlus/slvCtrlPlusDeviceFactory.ts @@ -9,7 +9,7 @@ import SlvCtrlProtocol, { DeviceInfo } from './slvCtrlProtocol.js'; import { getErrorFromDecodeResult } from '../deviceProtocol.js'; import EventEmitterFactory from '../../../factory/eventEmitterFactory.js'; import { SlvCtrlPlusDeviceAttributes } from './slvCtrlPlusDevice.js'; -import { DeviceId } from '../../deviceId.js'; +import { DeviceId, DetectionId } from '../../deviceId.js'; export default class SlvCtrlPlusDeviceFactory { @@ -33,10 +33,11 @@ export default class SlvCtrlPlusDeviceFactory this.logger = logger.child({ name: SlvCtrlPlusDeviceFactory.name }); } - public async create(deviceId: DeviceId, transport: DeviceBidirectionalTransport, provider: string): Promise { + public async create(detectionId: DetectionId, transport: DeviceBidirectionalTransport, provider: string): Promise { const deviceInfo = await this.getDeviceInfo(transport); const protocol = deviceInfo.protocol; - const knownDevice = this.knownDeviceRegistry.resolve(deviceId, deviceInfo.deviceType, provider); + + const knownDevice = this.knownDeviceRegistry.resolve(DeviceId.fromDetectionId(detectionId), deviceInfo.deviceType, provider); const deviceAttributes = await this.getAttributes(transport, protocol); const device = new GenericSlvCtrlPlusDevice( diff --git a/src/device/protocol/virtual/virtualDevice.ts b/src/device/protocol/virtual/virtualDevice.ts index 052948ea..74c2912f 100644 --- a/src/device/protocol/virtual/virtualDevice.ts +++ b/src/device/protocol/virtual/virtualDevice.ts @@ -20,8 +20,6 @@ export default class VirtualDevice< private readonly deviceLogic: TLogic; - private readonly logger: Logger; - private readonly statusUpdater?: NodeJS.Timeout; public constructor( @@ -36,12 +34,11 @@ export default class VirtualDevice< eventEmitter: EventEmitter, logger: Logger ) { - super(deviceId, deviceName, provider, connectedSince, false, deviceLogic.configureAttributes(), config, eventEmitter); + super(deviceId, deviceName, provider, connectedSince, false, deviceLogic.configureAttributes(), config, eventEmitter, logger); this.deviceModel = deviceModel; this.fwVersion = fwVersion; this.deviceLogic = deviceLogic; - this.logger = logger; } protected override async doClose(): Promise { diff --git a/src/device/protocol/virtual/virtualDeviceProvider.ts b/src/device/protocol/virtual/virtualDeviceProvider.ts index d6e7ece3..92e1bb83 100644 --- a/src/device/protocol/virtual/virtualDeviceProvider.ts +++ b/src/device/protocol/virtual/virtualDeviceProvider.ts @@ -10,6 +10,7 @@ import VirtualDeviceFactory from './virtualDeviceFactory.js'; import DeviceManager from '../../deviceManager.js'; import { asyncHandler } from '../../../util/async.js'; import { logError } from '../../../util/error.js'; +import { DetectionId } from '../../deviceId.js'; export type VirtualDeviceDetectionInfo = DeviceDetectionInfo & { type: 'virtual'; @@ -86,7 +87,7 @@ export default class VirtualDeviceProvider extends DeviceProvider; - private logger: Logger; - public constructor( deviceId: DeviceId, deviceName: string, @@ -84,13 +82,12 @@ export default class Zc95Device extends PeripheralDevice this.onReceivedMessage(data)); this.messageResponseHandler = messageResponseHandler; - this.logger = logger.child({ name: `${Zc95Device.name}.${transport.getDeviceIdentifier()}` }); } public async setAttribute< diff --git a/src/device/protocol/zc95/zc95DeviceFactory.ts b/src/device/protocol/zc95/zc95DeviceFactory.ts index 51b8b43d..bda9e140 100644 --- a/src/device/protocol/zc95/zc95DeviceFactory.ts +++ b/src/device/protocol/zc95/zc95DeviceFactory.ts @@ -12,7 +12,7 @@ import DeviceBidirectionalTransport from '../../transport/deviceBidirectionalTra import MessageResponseHandler from '../messageResponseHandler.js'; import EventEmitterFactory from '../../../factory/eventEmitterFactory.js'; import { logError } from '../../../util/error.js'; -import { DeviceId } from '../../deviceId.js'; +import { DeviceId, DetectionId } from '../../deviceId.js'; export default class Zc95DeviceFactory { @@ -37,7 +37,7 @@ export default class Zc95DeviceFactory } public async create( - deviceId: DeviceId, + detectionId: DetectionId, versionDetails: VersionMsgResponse, protocol: Zc95Protocol, transport: DeviceBidirectionalTransport, @@ -57,7 +57,9 @@ export default class Zc95DeviceFactory // We only receive serial no. info for ZC95 devices with fw >=2.0 const knownDevice = this.knownDeviceRegistry.resolve( - versionDetails.SerialNo !== undefined ? DeviceId.create(versionDetails.SerialNo) : deviceId, + versionDetails.SerialNo !== undefined + ? DeviceId.create(versionDetails.SerialNo) + : DeviceId.fromDetectionId(detectionId), 'zc95', provider, ); diff --git a/src/device/serializedTypes.ts b/src/device/serializedTypes.ts index 2c407fda..d4eebdcf 100644 --- a/src/device/serializedTypes.ts +++ b/src/device/serializedTypes.ts @@ -1,5 +1,6 @@ import DeviceState from './deviceState.js'; import { DeviceAttributeModifier } from './attribute/deviceAttribute.js'; +import { DeviceId } from './deviceId.js'; type SerializedDeviceAttributeBase = { name: string; @@ -55,7 +56,7 @@ export type SerializedDeviceAttribute = type SerializedDeviceBase = { connectedSince: Date; - deviceId: string; + deviceId: DeviceId; deviceName: string; provider: string; state: DeviceState; diff --git a/src/device/transport/bleObserver.ts b/src/device/transport/bleObserver.ts index e4cce4f3..9d5a4ab7 100644 --- a/src/device/transport/bleObserver.ts +++ b/src/device/transport/bleObserver.ts @@ -2,7 +2,7 @@ import noble, { Peripheral } from '@stoprocent/noble'; import Logger from '../../logging/Logger.js'; import DeviceManager, { DeviceDetectionInfo } from '../deviceManager.js'; import { logError } from '../../util/error.js'; -import { DeviceId } from '../deviceId.js'; +import { DetectionId } from '../deviceId.js'; import SharedObserver from './sharedObserver.js'; export type BleDeviceDetectionInfo = DeviceDetectionInfo & { @@ -72,7 +72,7 @@ export default class BleObserver extends SharedObserver const deviceInfo: BleDeviceDetectionInfo = { type: 'ble', - detectionId: DeviceId.create(peripheral.id), + detectionId: DetectionId.create(peripheral.id), peripheral, }; diff --git a/src/device/transport/serialPortObserver.ts b/src/device/transport/serialPortObserver.ts index 336d3591..61aac3f9 100644 --- a/src/device/transport/serialPortObserver.ts +++ b/src/device/transport/serialPortObserver.ts @@ -4,8 +4,10 @@ import Logger from '../../logging/Logger.js'; import DeviceManager, { DeviceDetectionInfo } from '../deviceManager.js'; import { usb } from 'usb'; import { logError } from '../../util/error.js'; -import { DeviceId } from '../deviceId.js'; +import { DetectionId } from '../deviceId.js'; import SharedObserver from './sharedObserver.js'; +import { CancellationToken, cancellationTokenReasons } from '@timesplinter/sequential-task-queue'; +import LatestOnlyTaskQueue from '../../util/latestOnlyTaskQueue.js'; export type SerialDeviceDetectionInfo = DeviceDetectionInfo & { type: 'serial'; @@ -22,7 +24,7 @@ export default class SerialPortObserver extends SharedObserver private rescanTimer?: NodeJS.Timeout; - private discoveryInFlight = false; + private readonly discoveryQueue: LatestOnlyTaskQueue = new LatestOnlyTaskQueue(); public constructor( deviceManager: DeviceManager, @@ -44,14 +46,13 @@ export default class SerialPortObserver extends SharedObserver } this.rescanTimer = setTimeout(() => { - if (this.discoveryInFlight) { - return; - } - this.discoveryInFlight = true; - this.discoverSerialDevices() - .catch(e => logError(this.logger, 'Error while scanning for new serial devices', e)) - .finally(() => { - this.discoveryInFlight = false; + this.discoveryQueue.run((cancellationToken) => this.discoverSerialDevices(cancellationToken)) + .catch((e: unknown) => { + if (e === cancellationTokenReasons.cancel) { + this.logger.debug('Serial device discovery run cancelled (superseded by a later run or stopped)'); + } else { + logError(this.logger, 'Error while scanning for new serial devices', e); + } }); }, 1000); }; @@ -60,50 +61,72 @@ export default class SerialPortObserver extends SharedObserver usb.addEventListener('disconnect', this.onUsbEventRef); } - public async discoverSerialDevices(): Promise + /** + * Gives a provider that just joined an already-running observer a chance at devices detected + * before it subscribed - discoverSerialDevices() itself only announces newly-discovered + * ports, so a device already sitting unclaimed here (e.g. because no provider wanted it yet) + * would otherwise never be offered to this provider. + */ + protected override async onSubsequentStart(): Promise + { + // announceDetectedDevice() itself is a no-op for a device that's already connected, so + // there's no need to filter those out here first. + for (const deviceInfo of this.managedDevices.values()) { + this.deviceManager.announceDetectedDevice(deviceInfo); + } + } + + public async discoverSerialDevices(cancellationToken?: CancellationToken): Promise { const foundDevices: Map = new Map(); + let ports: PortInfo[]; try { - const ports = await SerialPort.list(); + ports = await SerialPort.list(); + } catch (err) { + logError(this.logger, 'Could not list serial ports', err); + return; + } + + if (true === cancellationToken?.cancelled) { + // Superseded by a later discovery run, or the observer was stopped + return; + } - // Iterate through all serial ports and add them to the managed devices and try to connect - for (const portInfo of ports) { - if (undefined === portInfo.vendorId || undefined === portInfo.productId) { - continue; - } + // Iterate through all serial ports and add them to the managed devices and try to connect + for (const portInfo of ports) { + if (undefined === portInfo.vendorId || undefined === portInfo.productId) { + continue; + } - // If the serial number is not defined, create a "unique" one based on vendorId and productId - if (undefined === portInfo.serialNumber) { - portInfo.serialNumber = `serial-${portInfo.vendorId}-${portInfo.productId}-${portInfo.locationId}`; - } + // If the serial number is not defined, create a "unique" one based on vendorId and productId + if (undefined === portInfo.serialNumber) { + portInfo.serialNumber = `serial-${portInfo.vendorId}-${portInfo.productId}-${portInfo.locationId}`; + } - foundDevices.set(portInfo.serialNumber, null); + foundDevices.set(portInfo.serialNumber, null); - if (!this.managedDevices.has(portInfo.serialNumber)) { - const deviceInfo: SerialDeviceDetectionInfo = { - type: 'serial', - detectionId: DeviceId.create(portInfo.serialNumber), - portInfo - }; + if (!this.managedDevices.has(portInfo.serialNumber)) { + const deviceInfo: SerialDeviceDetectionInfo = { + type: 'serial', + detectionId: DetectionId.create(portInfo.serialNumber), + portInfo + }; - this.managedDevices.set(portInfo.serialNumber, deviceInfo); - this.logger.debug(`Managed devices: ${this.managedDevices.size}`); + this.managedDevices.set(portInfo.serialNumber, deviceInfo); + this.logger.debug(`Managed devices: ${this.managedDevices.size}`); - this.deviceManager.announceDetectedDevice(deviceInfo); - } + this.deviceManager.announceDetectedDevice(deviceInfo); } + } - // Remove devices that are no longer present - for (const [key, deviceInfo] of this.managedDevices) { - if (!foundDevices.has(key)) { - this.deviceManager.revokeDetectedDevice(deviceInfo); - this.managedDevices.delete(key); - this.logger.info(`Managed devices: ${this.managedDevices.size}`); - } + // Remove devices that are no longer present + for (const [key, deviceInfo] of this.managedDevices) { + if (!foundDevices.has(key)) { + this.deviceManager.revokeDetectedDevice(deviceInfo); + this.managedDevices.delete(key); + this.logger.info(`Managed devices: ${this.managedDevices.size}`); } - } catch (err) { - logError(this.logger, 'Could not list serial ports', err); } } @@ -118,5 +141,15 @@ export default class SerialPortObserver extends SharedObserver usb.removeEventListener('disconnect', this.onUsbEventRef); this.onUsbEventRef = undefined; } + + // Cancels a discovery run that's still awaiting SerialPort.list(), so it can't write + // stale results back into managedDevices below once it settles + await this.discoveryQueue.cancel(); + + for (const deviceInfo of this.managedDevices.values()) { + this.deviceManager.revokeDetectedDevice(deviceInfo); + } + + this.managedDevices.clear(); } } diff --git a/src/device/transport/sharedObserver.ts b/src/device/transport/sharedObserver.ts index 5705241f..490a2ef9 100644 --- a/src/device/transport/sharedObserver.ts +++ b/src/device/transport/sharedObserver.ts @@ -14,17 +14,15 @@ export default abstract class SharedObserver } public async start(): Promise { - // Not already running and nobody else currently starting it either - I'm the first + const joiningRunningObserver = this.activeUsers > 0; + + // Not already running and nobody else currently starting it either if (this.activeUsers === 0 && this.startupPromise === undefined) { this.startupPromise = this.onFirstStart(); } if (this.startupPromise !== undefined) { try { - // Only count as an active user once startup actually succeeded - a rejection - // here propagates out of start() without incrementing, for every concurrent - // caller awaiting the same promise, so a failed startup doesn't leave anyone - // thinking the observer is running await this.startupPromise; } finally { this.startupPromise = undefined; @@ -36,6 +34,10 @@ export default abstract class SharedObserver if (this.activeUsers > 1) { this.logger.debug(`Already running, now used by ${this.activeUsers} provider(s)`); } + if (joiningRunningObserver) { + // Reannounce all detected devices by this observer to a provider joining later + await this.onSubsequentStart(); + } } public async stop(): Promise { @@ -62,4 +64,13 @@ export default abstract class SharedObserver * Runs once, when the last remaining caller releases this observer. */ protected abstract onLastStop(): Promise; + + /** + * Runs for every caller that joins an already-running observer (i.e. every start() call + * except the first). No-op by default - only relevant to observers with multiple consumers + * per instance, where a newly-joining consumer may need to catch up on state it missed. + */ + protected async onSubsequentStart(): Promise { + return Promise.resolve(); + } } diff --git a/src/serial/synchronousSerialPort.ts b/src/serial/synchronousSerialPort.ts index 4db0ad65..f6b066fd 100644 --- a/src/serial/synchronousSerialPort.ts +++ b/src/serial/synchronousSerialPort.ts @@ -1,15 +1,17 @@ -import { Readable, Writable } from 'stream'; +import { Readable } from 'stream'; +import { SerialPortStream } from '@serialport/stream'; import { cancellationTokenReasons, SequentialTaskQueue, TaskOptions } from '@timesplinter/sequential-task-queue'; -import { PortInfo } from '@serialport/bindings-interface'; +import { BindingInterface, PortInfo } from '@serialport/bindings-interface'; import Logger from '../logging/Logger.js'; -import { asyncHandler } from '../util/async.js'; import { logError } from '../util/error.js'; +type CloseHandler = () => Promise; + export default class SynchronousSerialPort { private reader: Readable; - private writer: Writable; + private writer: SerialPortStream; private readonly portInfo: PortInfo; @@ -19,13 +21,24 @@ export default class SynchronousSerialPort private closed = false; - public constructor(portInfo: PortInfo, reader: Readable, writer: Writable, logger: Logger) { + private closePromise?: Promise; + + private readonly closeSubscribers: CloseHandler[] = []; + + private readonly handleStreamClose = (): void => { + this.close().catch((err: unknown) => logError(this.logger, 'Error closing serial port after stream close', err)); + }; + + public constructor(portInfo: PortInfo, reader: Readable, writer: SerialPortStream, logger: Logger) { this.portInfo = portInfo; this.reader = reader; this.writer = writer; this.queue = new SequentialTaskQueue(); this.logger = logger; this.queue.on('error', (error: unknown) => logError(this.logger, 'Error in queued task', error)); + + this.writer.on('close', this.handleStreamClose); + this.reader.on('close', this.handleStreamClose); } public async write(data: Buffer): Promise { @@ -38,36 +51,51 @@ export default class SynchronousSerialPort this.reader.on('data', dataProcessor); } - public onClose(callback: () => Promise): void { - const runClose = asyncHandler( - callback, - (e: unknown) => logError(this.logger, 'Error in writer/reader close handler', e) - ); - - const handleClose = (): void => { - this.closed = true; - this.queue.cancel(); - // Prevent second call and clean up listeners - this.writer.off('close', handleClose); - this.reader.off('close', handleClose); - runClose(); - }; - - this.writer.on('close', handleClose); - this.reader.on('close', handleClose); + public onClose(handler: CloseHandler): void { + this.closeSubscribers.push(handler); } public isOpen(): boolean { return !this.closed && this.writer.writable && this.reader.readable; } - public close(): void { + public async close(): Promise { + if (undefined !== this.closePromise) { + return this.closePromise; + } + + this.closePromise = this.doClose(); + + await this.closePromise; + } + + private async doClose(): Promise { this.closed = true; - this.queue.cancel(); - this.writer.end(() => { - this.writer.destroy(); - this.reader.destroy(); - }); + void this.queue.close(true); + this.writer.off('close', this.handleStreamClose); + this.reader.off('close', this.handleStreamClose); + + if (this.writer.isOpen) { + await new Promise((resolve) => { + this.writer.close((err) => { + if (err) { + logError(this.logger, `Error while closing serial port '${this.portInfo.path}'`, err); + } + resolve(); + }); + }); + } + + this.writer.destroy(); + this.reader.destroy(); + + for (const subscriber of this.closeSubscribers) { + try { + await subscriber(); + } catch (e: unknown) { + logError(this.logger, 'Error in writer/reader close handler', e); + } + } } public async writeAndExpect(data: Buffer, timeoutMs = 1000): Promise { @@ -113,7 +141,9 @@ export default class SynchronousSerialPort try { return await this.queue.push(wrappedPromise, options); - } catch (e) { + } catch (e: unknown) { + removeListeners(); + let reason = `task cancelled for unknown reason`; if (e === cancellationTokenReasons.timeout) { @@ -122,8 +152,6 @@ export default class SynchronousSerialPort reason = 'task deliberately cancelled'; } - removeListeners(); - throw new Error(reason, { cause: e }); } } diff --git a/src/serviceProvider/deviceServiceProvider.ts b/src/serviceProvider/deviceServiceProvider.ts index 43809239..122baa81 100644 --- a/src/serviceProvider/deviceServiceProvider.ts +++ b/src/serviceProvider/deviceServiceProvider.ts @@ -7,7 +7,6 @@ import { starWarsNouns } from '../util/dictionary.js'; import BufferedDeviceUpdater from '../device/updater/bufferedDeviceUpdater.js'; import GenericDeviceUpdater from '../device/genericDeviceUpdater.js'; import SerialDeviceTransportFactory from '../device/transport/serialDeviceTransportFactory.js'; -import { AnyDevice } from '../device/device.js'; import DeviceProviderManager from '../device/provider/deviceProviderManager.js'; import { AnyDeviceProvider } from '../device/provider/deviceProvider.js'; import SlvCtrlPlusSerialDeviceProvider from '../device/protocol/slvCtrlPlus/slvCtrlPlusSerialDeviceProvider.js'; @@ -41,7 +40,6 @@ import BleObserver from '../device/transport/bleObserver.js'; import AiroticDeviceProvider from '../device/protocol/airotic/airoticDeviceProvider.js'; import AiroticDeviceFactory from '../device/protocol/airotic/airoticDeviceFactory.js'; import DeviceProviderFactory from '../device/provider/deviceProviderFactory.js'; -import { DeviceId } from '../device/deviceId.js'; import KnownDeviceRegistry from '../device/knownDeviceRegistry.js'; export default class DeviceServiceProvider implements ServiceProvider { @@ -76,7 +74,6 @@ export default class DeviceServiceProvider implements ServiceProvider { return new DeviceManager( container.get('factory.eventEmitter').create(), - new Map(), container.get('settings.manager'), container.get('logger.default') ); diff --git a/src/serviceProvider/settingsServiceProvider.ts b/src/serviceProvider/settingsServiceProvider.ts index 1ddb3b03..c4d16282 100644 --- a/src/serviceProvider/settingsServiceProvider.ts +++ b/src/serviceProvider/settingsServiceProvider.ts @@ -22,17 +22,13 @@ export default class SettingsServiceProvider implements ServiceProvider { diff --git a/src/settings/settingsManager.ts b/src/settings/settingsManager.ts index 851e29ad..02cefc12 100644 --- a/src/settings/settingsManager.ts +++ b/src/settings/settingsManager.ts @@ -1,4 +1,5 @@ import fs from 'fs'; +import { watch, FSWatcher } from 'chokidar'; import PlainToClassSerializer from '../serialization/plainToClassSerializer.js'; import ClassToPlainSerializer from '../serialization/classToPlainSerializer.js'; import Settings, { SettingsSchema } from './settings.js'; @@ -6,7 +7,6 @@ import onChange from 'on-change'; import DeviceSource from './deviceSource.js'; import SlvCtrlPlusSerialDeviceProvider from '../device/protocol/slvCtrlPlus/slvCtrlPlusSerialDeviceProvider.js'; import Logger from '../logging/Logger.js'; -import SchemaValidationError from '../schemaValidation/schemaValidationError.js'; import EventEmitter from 'events'; import SettingsEventType from './settingsEventType.js'; import { JsonObject } from '../types.js'; @@ -31,6 +31,11 @@ export default class SettingsManager private readonly logger: Logger; + private watcher?: FSWatcher; + + /** Content of the settings file as written by our own save(), used to tell apart external edits from our own writes */ + private lastWrittenContent?: string; + public constructor( settingsFilePath: string, plainToClassSerializer: PlainToClassSerializer, @@ -54,20 +59,17 @@ export default class SettingsManager this.settings = SettingsManager.getDefaultSettings(); this.save(); } else { - const plainJsonSettings: JsonObject = JSON.parse(fs.readFileSync(this.settingsFilePath, 'utf8')); + const fileContent = fs.readFileSync(this.settingsFilePath, 'utf8'); try { - this.settings = this.plainToClassSerializer.transform(Settings, plainJsonSettings, SettingsSchema); + const plainJsonSettings: JsonObject = JSON.parse(fileContent); + this.settings = this.transformPlainToSettings(plainJsonSettings); } catch (e: unknown) { - if (!(e instanceof SchemaValidationError)) { - throw e; - } - - const invalidFormatMsg = `Settings are not in a valid format: ${e.message}`; - this.logger.error(invalidFormatMsg); - throw new Error(invalidFormatMsg, { cause: e }); + logError(this.logger, 'Settings are not in a valid format', e); + throw e; } + this.lastWrittenContent = fileContent; this.logger.info(`Settings loaded from file: ${this.settingsFilePath}`); } @@ -82,6 +84,40 @@ export default class SettingsManager this.logger.info(`Settings have been replaced with new value`); } + /** + * Starts watching the settings file on disk for changes made by another process (e.g. manual edits). + * Own writes via save() are recognized and ignored so they don't trigger a redundant reload. + */ + public startWatching(): void { + if (undefined !== this.watcher) { + return; + } + + this.watcher = watch(this.settingsFilePath, { + ignoreInitial: true, + // Debounce editors/tools that write the file in multiple chunks (or via temp file + rename) + awaitWriteFinish: { stabilityThreshold: 200, pollInterval: 50 }, + }); + + this.watcher.on('change', () => this.handleExternalChange()); + this.watcher.on('error', (err: unknown) => logError( + this.logger, + `Settings file watcher error for '${this.settingsFilePath}'`, + err + )); + + this.logger.debug(`Watching '${this.settingsFilePath}' for external changes`); + } + + public async stopWatching(): Promise { + if (undefined === this.watcher) { + return; + } + + await this.watcher.close(); + this.watcher = undefined; + } + public on (event: E, listener: SettingsEvents[E]): this { this.eventEmitter.on(event, listener); @@ -105,10 +141,12 @@ export default class SettingsManager try { const normalized = this.classToPlainSerializer.transform(this.settings); - fs.writeFileSync( - this.settingsFilePath, - JSON.stringify(normalized, null, 4) - ); + const json = JSON.stringify(normalized, null, 4); + + fs.writeFileSync(this.settingsFilePath, json); + // Remember what we just wrote so the file watcher can recognize and ignore this write + this.lastWrittenContent = json; + this.eventEmitter.emit('settingsChanged', this.settings); this.logger.debug(`Settings saved to '${this.settingsFilePath}' due to a change`); } catch (err: unknown) { @@ -116,6 +154,53 @@ export default class SettingsManager } } + /** + * Invoked by the file watcher when the settings file changed on disk. Ignores changes that + * originated from our own save() and reloads + emits an event for genuine external edits. + */ + private handleExternalChange(): void { + let content: string; + + try { + content = fs.readFileSync(this.settingsFilePath, 'utf8'); + } catch (err: unknown) { + logError(this.logger, `Could not read settings file '${this.settingsFilePath}' after change was detected`, err); + return; + } + + if (content === this.lastWrittenContent) { + // This change was caused by our own save(), nothing to do + return; + } + + let plainJsonSettings: JsonObject; + + try { + plainJsonSettings = JSON.parse(content); + } catch (err: unknown) { + logError(this.logger, `Ignoring external change to '${this.settingsFilePath}': content is not valid JSON`, err); + return; + } + + let parsedSettings: Settings; + + try { + parsedSettings = this.transformPlainToSettings(plainJsonSettings); + } catch (err: unknown) { + logError(this.logger, `Ignoring external change to '${this.settingsFilePath}': settings are not in a valid format`, err); + return; + } + + this.lastWrittenContent = content; + this.settings = onChange(parsedSettings, () => this.save()); + this.eventEmitter.emit('settingsChanged', this.settings); + this.logger.info(`Settings reloaded after external change to '${this.settingsFilePath}'`); + } + + private transformPlainToSettings(plainJsonSettings: JsonObject): Settings { + return this.plainToClassSerializer.transform(Settings, plainJsonSettings, SettingsSchema); + } + private static getDefaultSettings(): Settings { const settings = new Settings(); diff --git a/src/util/latestOnlyTaskQueue.ts b/src/util/latestOnlyTaskQueue.ts new file mode 100644 index 00000000..71c0c45f --- /dev/null +++ b/src/util/latestOnlyTaskQueue.ts @@ -0,0 +1,16 @@ +import { SequentialTaskQueue, CancellationToken, CancellablePromiseLike } from '@timesplinter/sequential-task-queue'; + +export default class LatestOnlyTaskQueue { + private queue = new SequentialTaskQueue(); + + public run( + task: (token: CancellationToken) => Promise + ): CancellablePromiseLike { + void this.queue.cancel(); + return this.queue.push(task); + } + + public cancel(reason?: unknown): Promise { + return this.queue.cancel(reason); + } +} \ No newline at end of file diff --git a/tests/integration/devices/sharedSerialObserver.spec.ts b/tests/integration/devices/sharedSerialObserver.spec.ts new file mode 100644 index 00000000..706f0910 --- /dev/null +++ b/tests/integration/devices/sharedSerialObserver.spec.ts @@ -0,0 +1,148 @@ +import { afterAll, assert, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'; +import { io as ioClient } from 'socket.io-client'; +import SlvCtrlPlusSerialDeviceProvider from '../../../src/device/protocol/slvCtrlPlus/slvCtrlPlusSerialDeviceProvider.js'; +import Zc95SerialDeviceProvider from '../../../src/device/protocol/zc95/zc95SerialDeviceProvider.js'; +import EStim2bSerialDeviceProvider from '../../../src/device/protocol/estim2b/estim2bSerialDeviceProvider.js'; +import WebSocketEvent from '../../../src/device/webSocketEvent.js'; +import { SlvCtrlPlusDeviceSimulator } from '../helpers/slvCtrlPlusDeviceSimulator.js'; +import { createTestApp, teardownTestApp, waitForNextWsEvent, createWsClient, TestApp } from '../helpers/appHelper.js'; +import Settings from '../../../src/settings/settings.js'; +import DeviceSource from '../../../src/settings/deviceSource.js'; + +process.env.LOG_LEVEL = process.env.LOG_LEVEL ?? 'silent'; + +const V1_PORT_PATH = '/dev/test-slvctrl-v1-0'; + +const SERIAL_SOURCE_ID = 'e6f7a8b9-6789-4321-abcd-ef1234567895'; + +// All three serial-protocol providers share the same SerialPortObserver singleton (they all +// watch the same physical ports and race to claim any detected device - see +// DeviceServiceProvider, where all three factories are wired to the same 'device.observer.serial' +// instance). This file exercises that shared-observer behavior specifically - it uses +// SlvCtrlPlusDeviceSimulator as a concrete stand-in device only because it's the most complete +// simulator available; the scenarios below are about the provider-manager/observer lifecycle, +// not SlvCtrl+ protocol semantics (those are covered in slvCtrlSerialDevice.spec.ts). +const ZC95_SOURCE_ID = 'a1c2d3e4-1234-4321-abcd-ef1234567896'; +const ESTIM2B_SOURCE_ID = 'b2d3e4f5-2345-4321-abcd-ef1234567897'; + +const serialSettings = { + knownDevices: {}, + deviceSources: { + // SlvCtrl+ listed first so it deterministically wins the claim race for the mock device + // used throughout this file (see DeviceManager.acquireDetectedDevice()'s FCFS queue). + [SERIAL_SOURCE_ID]: { + id: SERIAL_SOURCE_ID, + type: SlvCtrlPlusSerialDeviceProvider.providerName, + config: {}, + }, + [ZC95_SOURCE_ID]: { + id: ZC95_SOURCE_ID, + type: Zc95SerialDeviceProvider.providerName, + config: {}, + }, + [ESTIM2B_SOURCE_ID]: { + id: ESTIM2B_SOURCE_ID, + type: EStim2bSerialDeviceProvider.providerName, + config: {}, + }, + }, +}; + +describe('Serial protocol providers sharing one SerialPortObserver', () => { + let app: TestApp; + let wsEmitSpy: ReturnType; + let wsClient: ReturnType; + + beforeAll(async () => { + app = await createTestApp(serialSettings); + + wsEmitSpy = vi.spyOn(app.websocket, 'emit'); + + wsClient = await createWsClient(app.httpServer); + }); + + afterAll(async () => { + wsClient.disconnect(); + await teardownTestApp(app); + app.mockSerialPortFactory.reset(); + }); + + beforeEach(async () => { + // Close all connected devices before destroying the mock binding so their polling + // timers are stopped and the device-close chain completes cleanly. Without this, + // stale devices accumulate across test iterations: each one keeps a 100ms polling + // timer alive and floods the event loop with I/O errors after the binding is torn down. + await app.container.get('device.manager').reset(); + app.mockSerialPortFactory.reset(); + await app.container.get('device.observer.serial').discoverSerialDevices(); + wsEmitSpy.mockClear(); + }); + + it('emits deviceDisconnected only once per close(), not once per underlying stream', async () => { + // close() destroys the writer/reader after releasing the port, which triggers their own + // 'close' events - without a guard against that self-triggered re-entry, the device's + // close chain (and thus this WS event) would fire a second time for the same close(). + const simulator = new SlvCtrlPlusDeviceSimulator({ protocol: 'v1', deviceType: 'testDeviceV1' }); + app.mockSerialPortFactory.attachDevice(V1_PORT_PATH, simulator); + + const deviceConnected = waitForNextWsEvent(wsEmitSpy, WebSocketEvent.deviceConnected); + await app.container.get('device.observer.serial').discoverSerialDevices(); + const [payload] = await deviceConnected; + + const device = app.container.get('device.manager').getConnectedDevice(payload.deviceId); + assert(device !== null); + + wsEmitSpy.mockClear(); + + const deviceDisconnected = waitForNextWsEvent(wsEmitSpy, WebSocketEvent.deviceDisconnected); + await device.close(); + await deviceDisconnected; + + // Give any stray re-entrant close/emit a chance to fire before counting + await new Promise((resolve) => setTimeout(resolve, 200)); + + const disconnectCalls = wsEmitSpy.mock.calls.filter((call: unknown[]) => call[0] === WebSocketEvent.deviceDisconnected); + expect(disconnectCalls.length).toBe(1); + }); + + it('reconnects a device after its source is disabled then re-enabled, even while other sources keep the observer alive', async () => { + const simulator = new SlvCtrlPlusDeviceSimulator({ protocol: 'v1', deviceType: 'testDeviceV1' }); + app.mockSerialPortFactory.attachDevice(V1_PORT_PATH, simulator); + + const deviceConnected = waitForNextWsEvent(wsEmitSpy, WebSocketEvent.deviceConnected); + await app.container.get('device.observer.serial').discoverSerialDevices(); + await deviceConnected; + + expect(app.container.get('device.manager').getConnectedDevices().length).toBe(1); + + const settingsManager = app.container.get('settings.manager'); + + // Disable only the SlvCtrl+ source: the connected device must close and disappear. zc95 + // and estim2b stay enabled - they share the same SerialPortObserver, so this exercises + // the case where the observer stays alive (kept running by the other two sources) + // instead of fully stopping. + const disabledSettings = new Settings(); + disabledSettings.addDeviceSource(new DeviceSource(SERIAL_SOURCE_ID, SlvCtrlPlusSerialDeviceProvider.providerName, {}, false)); + disabledSettings.addDeviceSource(new DeviceSource(ZC95_SOURCE_ID, Zc95SerialDeviceProvider.providerName, {}, true)); + disabledSettings.addDeviceSource(new DeviceSource(ESTIM2B_SOURCE_ID, EStim2bSerialDeviceProvider.providerName, {}, true)); + + const deviceDisconnected = waitForNextWsEvent(wsEmitSpy, WebSocketEvent.deviceDisconnected); + settingsManager.replace(disabledSettings); + await deviceDisconnected; + + expect(app.container.get('device.manager').getConnectedDevices().length).toBe(0); + + // Re-enable the SlvCtrl+ source: the still-plugged-in device must reconnect on its own, + // even though the observer never fully stopped in between (zc95/estim2b kept it alive). + const reenabledSettings = new Settings(); + reenabledSettings.addDeviceSource(new DeviceSource(SERIAL_SOURCE_ID, SlvCtrlPlusSerialDeviceProvider.providerName, {}, true)); + reenabledSettings.addDeviceSource(new DeviceSource(ZC95_SOURCE_ID, Zc95SerialDeviceProvider.providerName, {}, true)); + reenabledSettings.addDeviceSource(new DeviceSource(ESTIM2B_SOURCE_ID, EStim2bSerialDeviceProvider.providerName, {}, true)); + + const deviceReconnected = waitForNextWsEvent(wsEmitSpy, WebSocketEvent.deviceConnected); + settingsManager.replace(reenabledSettings); + await deviceReconnected; + + expect(app.container.get('device.manager').getConnectedDevices().length).toBe(1); + }); +}); diff --git a/tests/integration/helpers/slvCtrlPlusDeviceSimulator.ts b/tests/integration/helpers/slvCtrlPlusDeviceSimulator.ts index 8dfee50a..e5160951 100644 --- a/tests/integration/helpers/slvCtrlPlusDeviceSimulator.ts +++ b/tests/integration/helpers/slvCtrlPlusDeviceSimulator.ts @@ -105,7 +105,12 @@ export class SlvCtrlPlusDeviceSimulator { // Use setImmediate so the response lands after writeAndExpect sets up its // data listener (the listener is registered before write() returns). setImmediate(() => { - bindingPort.emitData(Buffer.from(response + '\n', 'utf-8')); + // The port may have been closed by the time this fires (e.g. a disconnect + // racing this response) - emitData() throws if so, and nothing's listening + // for it anymore anyway. + if (bindingPort.isOpen) { + bindingPort.emitData(Buffer.from(response + '\n', 'utf-8')); + } }); } }; diff --git a/tests/unit/automation/scriptRuntime.spec.ts b/tests/unit/automation/scriptRuntime.spec.ts index 505f0610..80b5bc5a 100644 --- a/tests/unit/automation/scriptRuntime.spec.ts +++ b/tests/unit/automation/scriptRuntime.spec.ts @@ -21,6 +21,9 @@ class StubDevice extends Device { public readonly setAttributeCalls: Array<[string, AttributeValue]> = []; public constructor(id: DeviceId, name: string) { + const logger = mock(); + logger.child.mockReturnValue(mock()); + super( id, name, 'test', new Date(), true, { @@ -28,7 +31,7 @@ class StubDevice extends Device { 'label', undefined, DeviceAttributeModifier.readWrite, 'hello' ), }, - {}, new EventEmitter(), + {}, new EventEmitter(), logger, ); } diff --git a/tests/unit/device/detectedDeviceOfferQueue.spec.ts b/tests/unit/device/detectedDeviceOfferQueue.spec.ts index 10810af1..f7e6b202 100644 --- a/tests/unit/device/detectedDeviceOfferQueue.spec.ts +++ b/tests/unit/device/detectedDeviceOfferQueue.spec.ts @@ -5,14 +5,15 @@ import DetectedDeviceOfferQueue from '../../../src/device/detectedDeviceOfferQue import DeviceOfferRejectedError from '../../../src/device/deviceOfferRejectedError.js'; import { DeviceDetectionInfo } from '../../../src/device/deviceManager.js'; import { AnyDevice } from '../../../src/device/device.js'; -import { DeviceId } from '../../../src/device/deviceId.js'; +import { DeviceId, DetectionId } from '../../../src/device/deviceId.js'; import Logger from '../../../src/logging/Logger.js'; import TestDevice from './testDevice.js'; describe('DetectedDeviceOfferQueue', () => { let mockedLogger: ReturnType>; - const deviceId = DeviceId.create('device-1'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; + const detectionId = DetectionId.create('device-1'); + const deviceId = DeviceId.fromDetectionId(detectionId); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId }; beforeEach(() => { mockedLogger = mock(); @@ -23,11 +24,11 @@ describe('DetectedDeviceOfferQueue', () => { it('is false before any offer and true while one is pending', async () => { const queue = new DetectedDeviceOfferQueue(mockedLogger); - expect(queue.has(deviceId)).toBe(false); + expect(queue.has(detectionId)).toBe(false); const pendingPromise = queue.offer(deviceInfo, () => new Promise(() => {})); - await vi.waitFor(() => expect(queue.has(deviceId)).toBe(true)); + await vi.waitFor(() => expect(queue.has(detectionId)).toBe(true)); queue.closeAll(new DeviceOfferRejectedError('test cleanup')); await pendingPromise; @@ -95,10 +96,10 @@ describe('DetectedDeviceOfferQueue', () => { await queue.offer(deviceInfo, () => Promise.resolve(new DeviceOfferRejectedError('rejected'))); - expect(queue.has(deviceId)).toBe(false); + expect(queue.has(detectionId)).toBe(false); const pendingPromise = queue.offer(deviceInfo, () => new Promise(() => {})); - expect(queue.has(deviceId)).toBe(true); + expect(queue.has(detectionId)).toBe(true); queue.closeAll(new DeviceOfferRejectedError('test cleanup')); await pendingPromise; @@ -224,10 +225,10 @@ describe('DetectedDeviceOfferQueue', () => { it('closes every open queue', async () => { const queue = new DetectedDeviceOfferQueue(mockedLogger); - const otherDeviceId = DeviceId.create('device-2'); + const otherDetectionId = DetectionId.create('device-2'); const firstResultPromise = queue.offer(deviceInfo, () => new Promise(() => {})); - const secondResultPromise = queue.offer({ type: 'test', detectionId: otherDeviceId }, () => new Promise(() => {})); + const secondResultPromise = queue.offer({ type: 'test', detectionId: otherDetectionId }, () => new Promise(() => {})); queue.closeAll(new DeviceOfferRejectedError('reset')); @@ -235,8 +236,8 @@ describe('DetectedDeviceOfferQueue', () => { expect(firstResult.successful).toBe(false); expect(secondResult.successful).toBe(false); - expect(queue.has(deviceId)).toBe(false); - expect(queue.has(otherDeviceId)).toBe(false); + expect(queue.has(detectionId)).toBe(false); + expect(queue.has(otherDetectionId)).toBe(false); }); }); @@ -244,7 +245,7 @@ describe('DetectedDeviceOfferQueue', () => { 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')); + queue.revoke(detectionId, new DeviceOfferRejectedError('device disappeared')); const result = await queue.offer(deviceInfo, () => Promise.reject(new Error('should never run'))); @@ -263,7 +264,7 @@ describe('DetectedDeviceOfferQueue', () => { await vi.waitFor(() => expect(offerStarted).toBe(true)); - queue.revoke(deviceId, new DeviceOfferRejectedError('device disappeared')); + queue.revoke(detectionId, new DeviceOfferRejectedError('device disappeared')); const result = await resultPromise; expect(result.successful).toBe(false); @@ -280,8 +281,8 @@ describe('DetectedDeviceOfferQueue', () => { 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); + queue.revoke(detectionId, new DeviceOfferRejectedError('device disappeared')); + queue.dropIfRevoked(detectionId); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); const result = await queue.offer(deviceInfo, () => Promise.resolve(device)); @@ -292,8 +293,8 @@ describe('DetectedDeviceOfferQueue', () => { 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); + expect(() => queue.dropIfRevoked(detectionId)).not.toThrow(); + expect(queue.has(detectionId)).toBe(false); }); }); }); diff --git a/tests/unit/device/deviceManager.spec.ts b/tests/unit/device/deviceManager.spec.ts index 3406c289..44cd7cfa 100644 --- a/tests/unit/device/deviceManager.spec.ts +++ b/tests/unit/device/deviceManager.spec.ts @@ -3,10 +3,10 @@ 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, { AnyDevice } from "../../../src/device/device.js"; +import { 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"; +import { DeviceId, DetectionId } from "../../../src/device/deviceId.js"; import SettingsManager from "../../../src/settings/settingsManager.js"; import Settings from "../../../src/settings/settings.js"; import KnownDevice from "../../../src/settings/knownDevice.js"; @@ -53,11 +53,11 @@ describe('deviceManager', () => { const mockedLogger = mock(); mockedLogger.child.mockReturnValue(mockedLogger); - const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, new Map(), mockedSettingsManager, mockedLogger); + const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, mockedSettingsManager, mockedLogger); 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 deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; // New device connected expect(deviceManager.getConnectedDevices().length).toBe(0); @@ -77,10 +77,9 @@ describe('deviceManager', () => { it('it removes device from managed devices and emits event on disconnect', async () => { - 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 deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; const mockedDeviceManagerEventEmitter = mock(); mockedDeviceManagerEventEmitter.emit.mockReturnValue(true); @@ -88,7 +87,7 @@ describe('deviceManager', () => { const mockedLogger = mock(); mockedLogger.child.mockReturnValue(mockedLogger); - const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, connectedDevices, mockedSettingsManager, mockedLogger); + const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, mockedSettingsManager, mockedLogger); await connectDevice(deviceManager, mockedDeviceManagerEventEmitter, deviceInfo, device); mockClear(mockedDeviceManagerEventEmitter); @@ -103,10 +102,9 @@ describe('deviceManager', () => { it('it emits an event on device update', async () => { - 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 deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; const mockedDeviceManagerEventEmitter = mock(); mockedDeviceManagerEventEmitter.emit.mockReturnValue(true); @@ -114,7 +112,7 @@ describe('deviceManager', () => { const mockedLogger = mock(); mockedLogger.child.mockReturnValue(mockedLogger); - const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, connectedDevices, mockedSettingsManager, mockedLogger); + const deviceManager = new DeviceManager(mockedDeviceManagerEventEmitter, mockedSettingsManager, mockedLogger); await connectDevice(deviceManager, mockedDeviceManagerEventEmitter, deviceInfo, device); mockClear(mockedDeviceManagerEventEmitter); @@ -137,19 +135,23 @@ describe('deviceManager', () => { mockedLogger.child.mockReturnValue(mockedLogger); }); - it('returns the device when found by uuid', () => { - const uuid = 'known-device-uuid'; - const device = mock(); - const connectedDevices = new Map([[uuid, device]]); - const manager = new DeviceManager(mock(), connectedDevices, mockedSettingsManager, mockedLogger); + it('returns the device when found by uuid', async () => { + const mockedEventEmitter = mock(); + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); + + const deviceId = DeviceId.create('known-device-uuid'); + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; + + await connectDevice(manager, mockedEventEmitter, deviceInfo, device); - expect(manager.getConnectedDevice(uuid)).toBe(device); + expect(manager.getConnectedDevice(deviceId)).toBe(device); }); it('returns null when device is not found', () => { - const manager = new DeviceManager(mock(), new Map(), mockedSettingsManager, mockedLogger); + const manager = new DeviceManager(mock(), mockedSettingsManager, mockedLogger); - expect(manager.getConnectedDevice('unknown-uuid')).toBeNull(); + expect(manager.getConnectedDevice(DeviceId.create('unknown-uuid'))).toBeNull(); }); }); @@ -157,7 +159,7 @@ describe('deviceManager', () => { let mockedLogger: ReturnType>; let mockedEventEmitter: ReturnType>; const deviceId = DeviceId.create('device-1'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; beforeEach(() => { mockedLogger = mock(); @@ -167,7 +169,7 @@ describe('deviceManager', () => { it('emits deviceDetected event for a newly seen device', () => { mockedEventEmitter.emit.mockReturnValue(true); - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); @@ -175,7 +177,7 @@ describe('deviceManager', () => { }); 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); + const manager = new DeviceManager(mockedEventEmitter, 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 @@ -192,7 +194,7 @@ describe('deviceManager', () => { }); 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); + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); mockedEventEmitter.emit.mockReturnValue(true); manager.announceDetectedDevice(deviceInfo); @@ -211,9 +213,12 @@ describe('deviceManager', () => { 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); + it('does not emit event when device is already connected', async () => { + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); + + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + await connectDevice(manager, mockedEventEmitter, deviceInfo, device); + mockClear(mockedEventEmitter); manager.announceDetectedDevice(deviceInfo); @@ -222,7 +227,7 @@ describe('deviceManager', () => { 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); + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); @@ -238,7 +243,7 @@ describe('deviceManager', () => { // 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); + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); @@ -264,7 +269,7 @@ describe('deviceManager', () => { const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); - const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); manager.announceDetectedDevice(deviceInfo); @@ -276,7 +281,7 @@ describe('deviceManager', () => { let mockedLogger: ReturnType>; let mockedEventEmitter: ReturnType>; const deviceId = DeviceId.create('device-2'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; beforeEach(() => { mockedLogger = mock(); @@ -286,7 +291,7 @@ describe('deviceManager', () => { }); it('runs the first offer immediately and adds the device on success', async () => { - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); const result = await connectDevice(manager, mockedEventEmitter, deviceInfo, device); @@ -296,7 +301,7 @@ describe('deviceManager', () => { }); it('clears the queue and re-allows announcing after the only offer fails', async () => { - const manager = new DeviceManager(mockedEventEmitter, new Map(), mockedSettingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, mockedSettingsManager, mockedLogger); let resultPromise!: ReturnType; reactToDetection(mockedEventEmitter, () => { @@ -319,7 +324,7 @@ describe('deviceManager', () => { let mockedLogger: ReturnType>; let mockedEventEmitter: ReturnType>; const deviceId = DeviceId.create('device-4'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; beforeEach(() => { mockedLogger = mock(); @@ -335,7 +340,7 @@ describe('deviceManager', () => { const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); - const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); // 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 @@ -363,7 +368,7 @@ describe('deviceManager', () => { const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); - const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); let resolveOffer!: (device: AnyDevice) => void; const offerPromise = new Promise((resolve) => { resolveOffer = resolve; }); @@ -418,7 +423,7 @@ describe('deviceManager', () => { const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(new Settings()); - const manager = new DeviceManager(mock(), new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mock(), settingsManager, mockedLogger); expect(manager.isDeviceEnabled(DeviceId.create('unknown'))).toBe(true); }); @@ -427,7 +432,7 @@ describe('deviceManager', () => { const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(undefined); - const manager = new DeviceManager(mock(), new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mock(), settingsManager, mockedLogger); expect(manager.isDeviceEnabled(DeviceId.create('unknown'))).toBe(true); }); @@ -440,7 +445,7 @@ describe('deviceManager', () => { const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); - const manager = new DeviceManager(mock(), new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mock(), settingsManager, mockedLogger); expect(manager.isDeviceEnabled(deviceId)).toBe(false); }); @@ -462,7 +467,7 @@ describe('deviceManager', () => { // 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 detectionId = DetectionId.create('disabled-device-detection'); const canonicalId = DeviceId.create('disabled-device-canonical'); const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId }; @@ -472,8 +477,7 @@ describe('deviceManager', () => { const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); - const connectedDevices = new Map(); - const manager = new DeviceManager(mockedEventEmitter, connectedDevices, settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); const device = new TestDevice(canonicalId, 'Foo', new Date(), false, new EventEmitter()); @@ -486,14 +490,14 @@ describe('deviceManager', () => { 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 deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(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(mockedEventEmitter, new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); @@ -514,17 +518,16 @@ 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 deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; const enabledSettings = new Settings(); enabledSettings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(enabledSettings); - const connectedDevices = new Map(); const mockedEventEmitter = mock(); mockedEventEmitter.emit.mockReturnValue(true); - const manager = new DeviceManager(mockedEventEmitter, connectedDevices, settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); await connectDevice(manager, mockedEventEmitter, deviceInfo, device); @@ -539,19 +542,58 @@ describe('deviceManager', () => { expect(manager.getConnectedDevices()).toHaveLength(0); }); + it('re-announces a connected device once its known device is re-enabled after being disabled mid-session', async () => { + // Unlike the offer-rejected path (covered below), this device was already fully + // connected when its known device got disabled - its provider never stops/restarts + // in this scenario, so nothing but applySettingsChange() itself can trigger a retry. + const deviceId = DeviceId.create('device-disabled-while-connected'); + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; + + const settings = new Settings(); + settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); + + const settingsManager = mock(); + settingsManager.getSettings.mockReturnValue(settings); + + const mockedEventEmitter = mock(); + mockedEventEmitter.emit.mockReturnValue(true); + + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); + + const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); + await connectDevice(manager, mockedEventEmitter, deviceInfo, device); + expect(manager.getConnectedDevices()).toHaveLength(1); + + mockClear(mockedEventEmitter); + + settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, false)); + await manager.onSettingsChanged(); + + expect(manager.getConnectedDevices()).toHaveLength(0); + expect(mockedEventEmitter.emit).not.toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); + + mockClear(mockedEventEmitter); + + // Re-enable it - must be re-announced even though its provider kept running + // throughout and was never given another chance to redetect it. + settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); + await manager.onSettingsChanged(); + + expect(mockedEventEmitter.emit).toHaveBeenCalledWith(DeviceManagerEvent.deviceDetected, deviceInfo); + }); + 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 deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; const settings = new Settings(); settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, true)); const settingsManager = mock(); settingsManager.getSettings.mockReturnValue(settings); - const connectedDevices = new Map(); const mockedEventEmitter = mock(); mockedEventEmitter.emit.mockReturnValue(true); - const manager = new DeviceManager(mockedEventEmitter, connectedDevices, settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); await connectDevice(manager, mockedEventEmitter, deviceInfo, device); @@ -564,7 +606,7 @@ describe('deviceManager', () => { 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'); + const detectionId = DetectionId.create('device-pending-2-detected'); const canonicalId = DeviceId.create('device-pending-2-canonical'); const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId }; @@ -578,7 +620,7 @@ describe('deviceManager', () => { const mockedEventEmitter = mock(); mockedEventEmitter.emit.mockReturnValue(true); - const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); // Simulate a provider that connected a device via the detected-device pipeline whose // final id turns out to belong to a disabled device. @@ -604,7 +646,7 @@ describe('deviceManager', () => { it('does not re-announce a still-disabled pending device', async () => { const deviceId = DeviceId.create('device-pending-3'); - const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: deviceId }; + const deviceInfo: DeviceDetectionInfo = { type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }; const settings = new Settings(); settings.addKnownDevice(new KnownDevice(deviceId, 'Foo', 'test', 'test', {}, false)); @@ -615,7 +657,7 @@ describe('deviceManager', () => { const mockedEventEmitter = mock(); mockedEventEmitter.emit.mockReturnValue(true); - const manager = new DeviceManager(mockedEventEmitter, new Map(), settingsManager, mockedLogger); + const manager = new DeviceManager(mockedEventEmitter, settingsManager, mockedLogger); const device = new TestDevice(deviceId, 'Foo', new Date(), false, new EventEmitter()); await connectDevice(manager, mockedEventEmitter, deviceInfo, device); diff --git a/tests/unit/device/protocol/zc95/zc95DeviceFactory.spec.ts b/tests/unit/device/protocol/zc95/zc95DeviceFactory.spec.ts index 343b61b1..d0cb44de 100644 --- a/tests/unit/device/protocol/zc95/zc95DeviceFactory.spec.ts +++ b/tests/unit/device/protocol/zc95/zc95DeviceFactory.spec.ts @@ -15,7 +15,7 @@ import Zc95MessageFactory, { VersionMsgResponse, } from '../../../../../src/device/protocol/zc95/zc95MessageFactory.js'; import { MsgAndResponseIdentifier } from '../../../../../src/device/protocol/zc95/zc95Protocol.js'; -import { DeviceId } from '../../../../../src/device/deviceId.js'; +import { DeviceId, DetectionId } from '../../../../../src/device/deviceId.js'; describe('Zc95DeviceFactory', () => { let knownDeviceRegistry: MockProxy; @@ -35,7 +35,8 @@ describe('Zc95DeviceFactory', () => { Patterns: [{ Type: 'PatternDetail', Id: 0, Name: 'Pattern A' }], }; - const transportDeviceId = DeviceId.create('transport-device-id'); + const transportDetectionId = DetectionId.create('transport-device-id'); + const transportDeviceId = DeviceId.fromDetectionId(transportDetectionId); const provider = 'usb'; function baseVersionDetails(overrides: Partial = {}): VersionMsgResponse { @@ -75,7 +76,7 @@ describe('Zc95DeviceFactory', () => { const factory = createFactory(); const device = await factory.create( - transportDeviceId, + transportDetectionId, baseVersionDetails({ SerialNo: undefined }), mockProtocol, mockTransport, @@ -95,7 +96,7 @@ describe('Zc95DeviceFactory', () => { const expectedDeviceId = DeviceId.create('ZC95-SERIAL-123'); const device = await factory.create( - transportDeviceId, + transportDetectionId, baseVersionDetails({ SerialNo: 'ZC95-SERIAL-123' }), mockProtocol, mockTransport, @@ -128,7 +129,7 @@ describe('Zc95DeviceFactory', () => { const factory = createFactory(); const device = await factory.create( - transportDeviceId, + transportDetectionId, baseVersionDetails({ SerialNo: 'ZC95-SERIAL-123' }), mockProtocol, mockTransport, @@ -149,7 +150,7 @@ describe('Zc95DeviceFactory', () => { await expect( factory.create( - transportDeviceId, + transportDetectionId, baseVersionDetails({ SerialNo: undefined }), mockProtocol, mockTransport, diff --git a/tests/unit/device/provider/deviceProvider.spec.ts b/tests/unit/device/provider/deviceProvider.spec.ts index 813a57d2..df9c9eaf 100644 --- a/tests/unit/device/provider/deviceProvider.spec.ts +++ b/tests/unit/device/provider/deviceProvider.spec.ts @@ -8,7 +8,7 @@ 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 { DeviceId, DetectionId } from '../../../../src/device/deviceId.js'; import TestDevice from '../testDevice.js'; class TestProvider extends DeviceProvider @@ -25,7 +25,7 @@ class TestProvider extends DeviceProvider } protected createDevice(deviceDetectionInfo: DeviceDetectionInfo): Promise { - return Promise.resolve(new TestDevice(deviceDetectionInfo.detectionId, 'Foo', new Date(), false, new EventEmitter())); + return Promise.resolve(new TestDevice(DeviceId.fromDetectionId(deviceDetectionInfo.detectionId), 'Foo', new Date(), false, new EventEmitter())); } protected override async doStart(): Promise { @@ -50,7 +50,7 @@ class DetectingTestProvider extends DeviceProvider { - return Promise.resolve(new TestDevice(deviceDetectionInfo.detectionId, 'Foo', new Date(), false, new EventEmitter())); + return Promise.resolve(new TestDevice(DeviceId.fromDetectionId(deviceDetectionInfo.detectionId), 'Foo', new Date(), false, new EventEmitter())); } // Exposes the protected getConnectedDevice() so tests can check the provider's own @@ -156,11 +156,11 @@ describe('DeviceProvider', () => { const logger = mock(); logger.child.mockReturnValue(logger); - const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + const deviceManager = new DeviceManager(new EventEmitter(), settingsManager, logger); const provider = new DetectingTestProvider(deviceManager); await provider.start(); - deviceManager.announceDetectedDevice({ type: 'test', detectionId: DeviceId.create('device-a') }); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DetectionId.create('device-a') }); await vi.waitFor(() => expect(deviceManager.getConnectedDevices()).toHaveLength(1)); await provider.stop(); @@ -168,7 +168,7 @@ describe('DeviceProvider', () => { expect(deviceManager.getConnectedDevices()).toHaveLength(0); await provider.start(); - deviceManager.announceDetectedDevice({ type: 'test', detectionId: DeviceId.create('device-b') }); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DetectionId.create('device-b') }); // If `stopped` were not reset by start(), handleDeviceDetection() would abort this // detection immediately and the device would never be added. @@ -184,7 +184,7 @@ describe('DeviceProvider', () => { const logger = mock(); logger.child.mockReturnValue(logger); - const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + const deviceManager = new DeviceManager(new EventEmitter(), settingsManager, logger); const provider = new DetectingTestProvider(deviceManager); await provider.start(); @@ -197,7 +197,7 @@ describe('DeviceProvider', () => { } }); - deviceManager.announceDetectedDevice({ type: 'test', detectionId: deviceId }); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }); await vi.waitFor(() => expect(deviceManager.getConnectedDevices()).toHaveLength(1)); @@ -211,7 +211,7 @@ describe('DeviceProvider', () => { const logger = mock(); logger.child.mockReturnValue(logger); - const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + const deviceManager = new DeviceManager(new EventEmitter(), settingsManager, logger); let resolveCreateDevice!: (device: AnyDevice) => void; const createDevicePromise = new Promise((resolve) => { resolveCreateDevice = resolve; }); @@ -219,7 +219,7 @@ describe('DeviceProvider', () => { const provider = new SlowCreateDeviceProvider(deviceManager, createDevicePromise); await provider.start(); - deviceManager.announceDetectedDevice({ type: 'test', detectionId: DeviceId.create('device-stopped') }); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DetectionId.create('device-stopped') }); // Provider is stopped while createDevice() is still pending await provider.stop(); @@ -241,11 +241,11 @@ describe('DeviceProvider', () => { const logger = mock(); logger.child.mockReturnValue(logger); - const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + const deviceManager = new DeviceManager(new EventEmitter(), 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') }); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DetectionId.create('device-throw') }); await vi.waitFor(() => expect(provider.onConnectFailedCalls).toBe(1)); }); @@ -261,20 +261,20 @@ describe('DeviceProvider', () => { const logger = mock(); logger.child.mockReturnValue(logger); - const deviceManager = new DeviceManager(new EventEmitter(), new Map(), settingsManager, logger); + const deviceManager = new DeviceManager(new EventEmitter(), settingsManager, logger); let closeSpy: ReturnType | undefined; const provider = new TrackingTestProvider( deviceManager, (deviceDetectionInfo) => { - const device = new TestDevice(deviceDetectionInfo.detectionId, 'Foo', new Date(), false, new EventEmitter()); + const device = new TestDevice(DeviceId.fromDetectionId(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 }); + deviceManager.announceDetectedDevice({ type: 'test', detectionId: DetectionId.fromDeviceId(deviceId) }); // Disabled devices are closed internally by offerDevice() once rejected - wait for // that deterministically instead of a fixed sleep. diff --git a/tests/unit/device/testDevice.ts b/tests/unit/device/testDevice.ts index eebed15f..70077c99 100644 --- a/tests/unit/device/testDevice.ts +++ b/tests/unit/device/testDevice.ts @@ -1,6 +1,8 @@ import { EventEmitter } from "events"; +import { mock } from 'vitest-mock-extended'; import Device, {AttributeKeyOf, AttributeValueOf, DeviceAttributes} from "../../../src/device/device.js"; import { DeviceId } from '../../../src/device/deviceId.js'; +import Logger from '../../../src/logging/Logger.js'; export default class TestDevice extends Device { @@ -11,7 +13,10 @@ export default class TestDevice extends Device controllable: boolean, eventEmitter: EventEmitter, ) { - super(deviceId, deviceName, 'dummy', connectedSince, controllable, {}, {}, eventEmitter); + const logger = mock(); + logger.child.mockReturnValue(mock()); + + super(deviceId, deviceName, 'dummy', connectedSince, controllable, {}, {}, eventEmitter, logger); } public async setAttribute< diff --git a/tests/unit/device/transport/serialPortObserver.spec.ts b/tests/unit/device/transport/serialPortObserver.spec.ts index 934ba529..e197ccf6 100644 --- a/tests/unit/device/transport/serialPortObserver.spec.ts +++ b/tests/unit/device/transport/serialPortObserver.spec.ts @@ -5,8 +5,9 @@ import { usb } from 'usb'; import DeviceManager from '../../../../src/device/deviceManager.js'; import Logger from '../../../../src/logging/Logger.js'; import SerialPortObserver from '../../../../src/device/transport/serialPortObserver.js'; -import { DeviceId } from '../../../../src/device/deviceId.js'; +import { DetectionId } from '../../../../src/device/deviceId.js'; import { waitTicks } from '../../helper/async.js'; +import { CancellationToken } from '@timesplinter/sequential-task-queue'; // usb is a real, module-wide EventTarget - without mocking it, addEventListener() calls made in // one test would still be registered when the next test runs, eventually tripping Node's @@ -111,7 +112,7 @@ describe('SerialPortObserver', () => { expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledOnce(); expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledWith( - expect.objectContaining({ detectionId: DeviceId.create('SN001'), portInfo: port }), + expect.objectContaining({ detectionId: DetectionId.create('SN001'), portInfo: port }), ); }); @@ -124,7 +125,7 @@ describe('SerialPortObserver', () => { const expectedSn = 'serial-0403-6001-port1'; expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledWith( - expect.objectContaining({ detectionId: DeviceId.create(expectedSn) }), + expect.objectContaining({ detectionId: DetectionId.create(expectedSn) }), ); }); @@ -151,7 +152,7 @@ describe('SerialPortObserver', () => { expect(mockDeviceManager.revokeDetectedDevice).toHaveBeenCalledOnce(); expect(mockDeviceManager.revokeDetectedDevice).toHaveBeenCalledWith( - expect.objectContaining({ detectionId: DeviceId.create('SN001') }), + expect.objectContaining({ detectionId: DetectionId.create('SN001') }), ); }); @@ -189,10 +190,10 @@ describe('SerialPortObserver', () => { expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledTimes(2); expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledWith( - expect.objectContaining({ detectionId: DeviceId.create('SN001') }), + expect.objectContaining({ detectionId: DetectionId.create('SN001') }), ); expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledWith( - expect.objectContaining({ detectionId: DeviceId.create('SN002') }), + expect.objectContaining({ detectionId: DetectionId.create('SN002') }), ); }); }); @@ -279,4 +280,82 @@ describe('SerialPortObserver', () => { await expect(observer.stop()).resolves.not.toThrow(); }); }); + + describe('restart after full stop (e.g. device source disabled then re-enabled)', () => { + it('re-announces a still-plugged-in device after the observer was fully stopped and started again', async () => { + const port = makePortInfo({ path: '/dev/ttyUSB0', serialNumber: 'SN001', vendorId: '0403', productId: '6001' }); + vi.spyOn(SerialPort, 'list').mockResolvedValue([port]); + const observer = createObserver(); + + await observer.start(); + expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledOnce(); + + await observer.stop(); // activeUsers 1 -> 0, triggers onLastStop() + await observer.start(); // brand new provider re-acquiring the shared observer + + expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledTimes(2); + }); + + it('revokes every still-tracked device via the device manager when fully stopped', async () => { + const port1 = makePortInfo({ path: '/dev/ttyUSB0', serialNumber: 'SN001', vendorId: '0403', productId: '6001' }); + const port2 = makePortInfo({ path: '/dev/ttyUSB1', serialNumber: 'SN002', vendorId: '0403', productId: '6015' }); + vi.spyOn(SerialPort, 'list').mockResolvedValue([port1, port2]); + const observer = createObserver(); + + await observer.start(); + await observer.stop(); + + expect(mockDeviceManager.revokeDetectedDevice).toHaveBeenCalledTimes(2); + expect(mockDeviceManager.revokeDetectedDevice).toHaveBeenCalledWith( + expect.objectContaining({ detectionId: DetectionId.create('SN001') }), + ); + expect(mockDeviceManager.revokeDetectedDevice).toHaveBeenCalledWith( + expect.objectContaining({ detectionId: DetectionId.create('SN002') }), + ); + }); + }); + + describe('discoverSerialDevices() cancellation', () => { + it('does not write results back into managedDevices when cancelled while awaiting SerialPort.list()', async () => { + const port = makePortInfo({ path: '/dev/ttyUSB0', serialNumber: 'SN001', vendorId: '0403', productId: '6001' }); + + let resolveList: (ports: PortInfoLike[]) => void = () => undefined; + vi.spyOn(SerialPort, 'list').mockImplementation(() => new Promise((resolve) => { resolveList = resolve; })); + + const observer = createObserver(); + const cancellationToken: CancellationToken = { cancelled: false, cancel: () => undefined }; + + const discoveryPromise = observer.discoverSerialDevices(cancellationToken); + + // Simulates onLastStop() cancelling this run (via discoveryQueue.cancel()) because + // the observer was stopped while this run was still in flight. + cancellationToken.cancelled = true; + + resolveList([port]); + await discoveryPromise; + + expect(mockDeviceManager.announceDetectedDevice).not.toHaveBeenCalled(); + expect(mockDeviceManager.revokeDetectedDevice).not.toHaveBeenCalled(); + }); + }); + + describe('catch-up for a provider joining an already-running observer', () => { + it('re-announces an already-managed device when a second provider starts', async () => { + const port = makePortInfo({ path: '/dev/ttyUSB0', serialNumber: 'SN001', vendorId: '0403', productId: '6001' }); + vi.spyOn(SerialPort, 'list').mockResolvedValue([port]); + const observer = createObserver(); + + await observer.start(); // first provider - full discovery, announces once + expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledOnce(); + + await observer.start(); // second provider joins while already running, without a rescan + + // Re-announced unconditionally - announceDetectedDevice() itself is a no-op for a + // device that's already claimed/connected, so the observer doesn't need to check first. + expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledTimes(2); + expect(mockDeviceManager.announceDetectedDevice).toHaveBeenCalledWith( + expect.objectContaining({ detectionId: DetectionId.create('SN001') }), + ); + }); + }); }); diff --git a/tests/unit/device/transport/sharedObserver.spec.ts b/tests/unit/device/transport/sharedObserver.spec.ts index 50f54ea2..09825d43 100644 --- a/tests/unit/device/transport/sharedObserver.spec.ts +++ b/tests/unit/device/transport/sharedObserver.spec.ts @@ -7,15 +7,25 @@ class TestObserver extends SharedObserver { public firstStarts = 0; public lastStops = 0; + public subsequentStarts = 0; public failNextStart = false; + // Allows tests to control when onFirstStart() resolves, to simulate a slow/in-flight startup. + private firstStartGate: Promise = Promise.resolve(); + public constructor() { super(mock()); } + public setFirstStartGate(gate: Promise): void { + this.firstStartGate = gate; + } + protected async onFirstStart(): Promise { this.firstStarts++; + await this.firstStartGate; + if (this.failNextStart) { this.failNextStart = false; throw new Error('startup failed'); @@ -25,6 +35,10 @@ class TestObserver extends SharedObserver protected async onLastStop(): Promise { this.lastStops++; } + + protected override async onSubsequentStart(): Promise { + this.subsequentStarts++; + } } describe('SharedObserver', () => { @@ -56,4 +70,60 @@ describe('SharedObserver', () => { await observer.stop(); expect(observer.lastStops).toBe(1); }); + + it('does not run onSubsequentStart() for the very first caller', async () => { + const observer = new TestObserver(); + + await observer.start(); + + expect(observer.subsequentStarts).toBe(0); + }); + + it('runs onSubsequentStart() for a caller joining an already-running observer', async () => { + const observer = new TestObserver(); + + await observer.start(); + await observer.start(); + + expect(observer.firstStarts).toBe(1); + expect(observer.subsequentStarts).toBe(1); + + await observer.start(); + expect(observer.subsequentStarts).toBe(2); + }); + + it('does not run onSubsequentStart() for concurrent first-time joiners', async () => { + const observer = new TestObserver(); + + let releaseFirstStart: () => void = () => undefined; + observer.setFirstStartGate(new Promise((resolve) => { releaseFirstStart = resolve; })); + + const firstStart = observer.start(); + const secondStart = observer.start(); + + releaseFirstStart(); + await Promise.all([firstStart, secondStart]); + + expect(observer.firstStarts).toBe(1); + expect(observer.subsequentStarts).toBe(0); + }); + + it('runs onSubsequentStart() again after a full stop and restart', async () => { + const observer = new TestObserver(); + + await observer.start(); + await observer.start(); + expect(observer.subsequentStarts).toBe(1); + + await observer.stop(); + await observer.stop(); + expect(observer.lastStops).toBe(1); + + await observer.start(); + expect(observer.firstStarts).toBe(2); + expect(observer.subsequentStarts).toBe(1); + + await observer.start(); + expect(observer.subsequentStarts).toBe(2); + }); }); diff --git a/tests/unit/settings/settingsManager.spec.ts b/tests/unit/settings/settingsManager.spec.ts new file mode 100644 index 00000000..8246dc57 --- /dev/null +++ b/tests/unit/settings/settingsManager.spec.ts @@ -0,0 +1,152 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { mock } from 'vitest-mock-extended'; +import { EventEmitter } from 'events'; +import fs from 'fs'; +import os from 'os'; +import path from 'path'; +import { Ajv2020 } from 'ajv/dist/2020.js'; +import ajvFormatsPlugin from 'ajv-formats'; +import SettingsManager from '../../../src/settings/settingsManager.js'; +import PlainToClassSerializer from '../../../src/serialization/plainToClassSerializer.js'; +import ClassToPlainSerializer from '../../../src/serialization/classToPlainSerializer.js'; +import Logger from '../../../src/logging/Logger.js'; +import DeviceSource from '../../../src/settings/deviceSource.js'; +import SettingsEventType from '../../../src/settings/settingsEventType.js'; + +const createSerializers = (): { plainToClass: PlainToClassSerializer; classToPlain: ClassToPlainSerializer } => { + const ajv = new Ajv2020({ allErrors: true, strict: true }); + ajvFormatsPlugin.default(ajv); + + return { + plainToClass: new PlainToClassSerializer(ajv, { excludeExtraneousValues: true }), + classToPlain: new ClassToPlainSerializer(), + }; +}; + +const waitFor = async (assertion: () => void, timeout = 3000, interval = 50): Promise => { + const start = Date.now(); + + for (;;) { + try { + assertion(); + return; + } catch (err) { + if (Date.now() - start > timeout) { + throw err; + } + await new Promise(resolve => setTimeout(resolve, interval)); + } + } +}; + +describe('SettingsManager', () => { + let tmpDir: string; + let settingsFilePath: string; + let mockLogger: ReturnType>; + let settingsManager: SettingsManager; + + beforeEach(() => { + tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'slvctrlplus-settingsmanager-test-')); + settingsFilePath = path.join(tmpDir, 'settings.json'); + fs.writeFileSync(settingsFilePath, JSON.stringify({ knownDevices: {}, deviceSources: {} })); + + mockLogger = mock(); + mockLogger.child.mockReturnValue(mockLogger); + + const { plainToClass, classToPlain } = createSerializers(); + + settingsManager = new SettingsManager( + settingsFilePath, + plainToClass, + classToPlain, + new EventEmitter(), + mockLogger, + ); + }); + + afterEach(async () => { + await settingsManager.stopWatching(); + fs.rmSync(tmpDir, { recursive: true, force: true }); + }); + + it('does not treat its own save() write as an external change', async () => { + const settings = settingsManager.load(); + settingsManager.startWatching(); + + const changedListener = vi.fn(); + settingsManager.on(SettingsEventType.changed, changedListener); + + settings.addDeviceSource(new DeviceSource('source-1', 'virtual', {})); + + // Own write should fire exactly once, from save() itself, never from the watcher + await new Promise(resolve => setTimeout(resolve, 500)); + + expect(changedListener).toHaveBeenCalledTimes(1); + expect(mockLogger.info).not.toHaveBeenCalledWith(expect.stringContaining('reloaded after external change')); + }); + + it('reloads settings and emits an event when the file is changed externally', async () => { + settingsManager.load(); + settingsManager.startWatching(); + await new Promise(resolve => setTimeout(resolve, 300)); // let chokidar finish its async setup + + const changedListener = vi.fn(); + settingsManager.on(SettingsEventType.changed, changedListener); + + const externalSourceId = 'b6a0f45e-c3d0-4dca-ab81-7daac0764292'; + const externalContent = JSON.stringify({ + knownDevices: {}, + deviceSources: { + [externalSourceId]: { id: externalSourceId, type: 'virtual', config: {} }, + }, + }); + + fs.writeFileSync(settingsFilePath, externalContent); + + await waitFor(() => expect(changedListener).toHaveBeenCalledTimes(1)); + + const reloaded = settingsManager.getSettings(); + expect(reloaded?.getDeviceSources().has(externalSourceId)).toBe(true); + }); + + it('ignores an externally written file with invalid JSON and keeps previous settings', async () => { + settingsManager.load(); + settingsManager.startWatching(); + await new Promise(resolve => setTimeout(resolve, 300)); // let chokidar finish its async setup + + const changedListener = vi.fn(); + settingsManager.on(SettingsEventType.changed, changedListener); + + fs.writeFileSync(settingsFilePath, '{ not valid json'); + + await waitFor(() => expect(mockLogger.error).toHaveBeenCalled()); + + expect(changedListener).not.toHaveBeenCalled(); + expect(settingsManager.getSettings()?.getDeviceSources().size).toBe(0); + }); + + it('ignores an externally written file that fails schema validation and keeps previous settings', async () => { + settingsManager.load(); + settingsManager.startWatching(); + await new Promise(resolve => setTimeout(resolve, 300)); // let chokidar finish its async setup + + const changedListener = vi.fn(); + settingsManager.on(SettingsEventType.changed, changedListener); + + fs.writeFileSync(settingsFilePath, JSON.stringify({ unexpectedField: true })); + + await waitFor(() => expect(mockLogger.error).toHaveBeenCalled()); + + expect(changedListener).not.toHaveBeenCalled(); + }); + + it('startWatching() is idempotent and stopWatching() can be called when not watching', async () => { + settingsManager.load(); + + settingsManager.startWatching(); + settingsManager.startWatching(); + + await settingsManager.stopWatching(); + await settingsManager.stopWatching(); + }); +});