Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
eb60135
wip
heavyrubberslave Jul 27, 2026
d23a8ed
fix: repair device offer queue hand-off, stop race and disabled-devic…
heavyrubberslave Jul 28, 2026
777a9b9
refactor: convert runNextInQueue to async/await
heavyrubberslave Jul 28, 2026
aba609f
refactor: clean up deviceManager/deviceProvider leftovers
heavyrubberslave Jul 28, 2026
d3a9448
Remove unused imports
heavyrubberslave Jul 28, 2026
fad4703
fix: reject stale offers that settle after a revoke/reset (coderabbit)
heavyrubberslave Jul 28, 2026
842582e
refactor: improve naming clarity in deviceManager.ts
heavyrubberslave Jul 28, 2026
cbf16ae
refactor: further naming clarity in deviceManager.ts
heavyrubberslave Jul 28, 2026
f583d16
Simplify runNextOfferInQueue, eliminate undefined from device offer c…
heavyrubberslave Jul 28, 2026
2b7e237
Extract detected-device offer queue mechanics into DetectedDeviceOffe…
heavyrubberslave Jul 29, 2026
9f5def9
Remove injected acceptor from DetectedDeviceOfferQueue
heavyrubberslave Jul 29, 2026
fd66dc9
Make DetectedDeviceOfferQueue.open() a no-op if already open
heavyrubberslave Jul 29, 2026
db2b598
Fix cancellation check and extract runOffer() in DetectedDeviceOfferQ…
heavyrubberslave Jul 29, 2026
84ca43b
Extract createAndRegisterDevice() from DeviceProvider.handleDeviceDet…
heavyrubberslave Jul 29, 2026
627183b
Address CodeRabbit PR #99 review threads
heavyrubberslave Jul 29, 2026
040eceb
Fix TOCTOU race parking a revoked disabled device for retry
heavyrubberslave Jul 29, 2026
6b51a0a
Fix CodeRabbit thread 2: discard a queue nobody actually offered to
heavyrubberslave Jul 30, 2026
22f5555
Simplify DetectedDeviceOfferQueue to lazy, self-opening queues
heavyrubberslave Aug 1, 2026
69e18f1
Clean up device offer rejection log messages
heavyrubberslave Aug 1, 2026
bc9414f
Fix dropIfRevoked() being unreachable behind the has() reentrancy guard
heavyrubberslave Aug 1, 2026
84964fe
Unify DetectedDeviceOfferQueue on close(), address remaining CodeRabb…
heavyrubberslave Aug 1, 2026
6ab565c
Better error handling for failed handshake
heavyrubberslave Aug 1, 2026
059d8fe
Cancel pending offers before closing connected devices in reset()
heavyrubberslave Aug 1, 2026
0bf02dc
Remove comment
heavyrubberslave Aug 1, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
"read-last-lines": "^1.8.0",
"reflect-metadata": "^0.1.13",
"say": "^0.16.0",
"sequential-task-queue": "^1.2.1",
"@timesplinter/sequential-task-queue": "^1.3.1",
"serialport": "^13.0.0",
"socket.io": "^4.8.3",
"speaker": "https://github.com/SlvCtrlPlus/node-speaker/releases/download/v0.1.0/speaker-v0.1.0.tgz",
Expand Down
152 changes: 152 additions & 0 deletions src/device/detectedDeviceOfferQueue.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
import { CancellationToken, sequentialTaskQueueEvents, SequentialTaskQueue } from '@timesplinter/sequential-task-queue';
import { AnyDevice } from './device.js';
import { DeviceDetectionInfo } from './deviceManager.js';
import DeviceOfferRejectedError from './deviceOfferRejectedError.js';
import Logger from '../logging/Logger.js';
import { logError } from '../util/error.js';

export type OfferResult<D extends AnyDevice> =
| { successful: true, device: D }
| { successful: false, reason: unknown };

type DeviceOffer<D extends AnyDevice> = (cancellationToken: CancellationToken) => Promise<D | DeviceOfferRejectedError>;

export default class DetectedDeviceOfferQueue
{
private readonly queues: Map<string, SequentialTaskQueue> = new Map();

private readonly logger: Logger;

public constructor(logger: Logger) {
this.logger = logger;
}

private getOrCreateQueue(detectionId: string): SequentialTaskQueue
{
let queue = this.queues.get(detectionId);

if (queue !== undefined) {
return queue;
}

queue = new SequentialTaskQueue();

queue.on(sequentialTaskQueueEvents.drained, () => {
// A revoked (closed) queue must survive its own drain - it's kept around deliberately
// as a tombstone so a late offer can still see it and reject itself.
if (!queue.isClosed) {
this.queues.delete(detectionId);
}
});

this.queues.set(detectionId, queue);

return queue;
}

public offer<D extends AnyDevice>(deviceDetectionInfo: DeviceDetectionInfo, deviceOffer: DeviceOffer<D>): Promise<OfferResult<D>>
{
const detectionId = deviceDetectionInfo.detectionId;
const queue = this.getOrCreateQueue(detectionId);

if (queue.isClosed) {
return Promise.resolve({
successful: false,
reason: new DeviceOfferRejectedError('Device is not available anymore for offering'),
});
}

const task = queue.push((cancellationToken: CancellationToken) => this.runOffer(deviceOffer, cancellationToken));

return Promise.resolve(task.then(
(result: OfferResult<D>): OfferResult<D> => {
if (result.successful) {
// Reject every other still-queued offer for this detection id without them
// ever running, since this device has already been claimed.
this.close(detectionId, new DeviceOfferRejectedError('Device has been claimed by another provider'));
}

return result;
},
// Reached either if the offer was cancelled while still queued (never even starting -
// our callback above never ran) or if deviceOffer() itself rejected/threw uncaught -
// translate both the same way.
(reason: unknown): OfferResult<D> => ({
successful: false,
reason: reason,
})
));
}

private async runOffer<D extends AnyDevice>(
deviceOffer: DeviceOffer<D>,
cancellationToken: CancellationToken
): Promise<OfferResult<D>> {
const device = await deviceOffer(cancellationToken);

if (device instanceof DeviceOfferRejectedError) {
return { successful: false, reason: device };
}

if (true !== cancellationToken.cancelled) {
return { successful: true, device };
}

// In case this offer lost the race against another offer: close the device and reject the offer with a meaningful reason.
try {
await device.close();
} catch (e: unknown) {
logError(this.logger, `Failed to close device '${device.getDeviceId}' offered after its queue was cleared`, e);
}

return {
successful: false,
reason: cancellationToken.reason,
};
}

/**
* True if a queue currently exists for this detection id - either genuinely active/in-flight,
* or a closed tombstone left behind by revoke(). Does not distinguish between the two;
* callers that need "is a fresh announce still blocked by a past revoke" must call
* dropIfRevoked() first.
*/
public has(detectionId: string): boolean
{
return this.queues.has(detectionId);
}

public dropIfRevoked(detectionId: string): void
{
const queue = this.queues.get(detectionId);

if (queue !== undefined && queue.isClosed) {
this.queues.delete(detectionId);
}
}

private close(detectionId: string, reason: DeviceOfferRejectedError): void
{
const queue = this.queues.get(detectionId);

if (undefined !== queue) {
void queue.close(true, reason);
}

this.queues.delete(detectionId);
}

public revoke(detectionId: string, reason: DeviceOfferRejectedError): void
{
const queue = this.getOrCreateQueue(detectionId);

void queue.close(true, reason);
}
Comment thread
heavyrubberslave marked this conversation as resolved.

public closeAll(reason: DeviceOfferRejectedError): void
{
for (const detectionId of this.queues.keys()) {
this.close(detectionId, reason);
}
}
}
Loading
Loading