Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
19 changes: 17 additions & 2 deletions lib/internal/vfs/watcher.js
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ const {
ArrayPrototypePush,
ObjectAssign,
Promise,
PromiseReject,
PromiseResolve,
SafeMap,
SafeSet,
Expand Down Expand Up @@ -579,6 +580,7 @@ const kMaxPendingEvents = 1024;
class VFSWatchAsyncIterable {
#watcher;
#closed = false;
#abortError = null;
#pendingEvents = [];
#pendingResolvers = [];

Expand Down Expand Up @@ -616,14 +618,21 @@ class VFSWatchAsyncIterable {
}
});

// Handle abort signal - reject pending next() with AbortError
// Handle abort signal - reject the next pending next() with AbortError,
// then behave like an exhausted iterator, matching native fs.promises.watch.
if (signal) {
const onAbort = () => {
this.#closed = true;
const err = new AbortError(undefined, { cause: signal.reason });
while (this.#pendingResolvers.length > 0) {
if (this.#pendingResolvers.length > 0) {
const { reject } = this.#pendingResolvers.shift();
reject(err);
while (this.#pendingResolvers.length > 0) {
const { resolve } = this.#pendingResolvers.shift();
resolve({ done: true, value: undefined });
}
} else {
this.#abortError = err;
}
this.#watcher.close();
};
Expand All @@ -648,6 +657,12 @@ class VFSWatchAsyncIterable {
* @returns {Promise<IteratorResult>}
*/
next() {
if (this.#abortError !== null) {
const err = this.#abortError;
this.#abortError = null;
return PromiseReject(err);
}

if (this.#closed) {
return PromiseResolve({ done: true, value: undefined });
}
Expand Down
4 changes: 3 additions & 1 deletion test/parallel/test-vfs-watch-abort-signal.js
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,13 @@ const vfs = require('node:vfs');
setImmediate(() => myVfs.writeFileSync('/file.txt', 'b'));
}

// promises.watch with pre-aborted signal resolves done immediately
// promises.watch with pre-aborted signal rejects the first next() with
// AbortError, matching native fs.promises.watch
{
const myVfs = vfs.create();
myVfs.writeFileSync('/p.txt', 'a');
const iter = myVfs.promises.watch('/p.txt', { signal: AbortSignal.abort() });
await assert.rejects(iter.next(), { name: 'AbortError', code: 'ABORT_ERR' });
const r = await iter.next();
assert.strictEqual(r.done, true);
}
Expand Down
Loading