@@ -4,6 +4,127 @@ import type * as opts from './types/options';
44import type * as misc from './types/misc' ;
55import type { FsCallbackApi , FsPromisesApi } from './types' ;
66
7+ // AsyncIterator implementation for promises.watch
8+ class FSWatchAsyncIterator implements AsyncIterableIterator < { eventType : string ; filename : string | Buffer } > {
9+ private watcher : any ;
10+ private eventQueue : Array < { eventType : string ; filename : string | Buffer } > = [ ] ;
11+ private resolveQueue : Array < { resolve : Function ; reject : Function } > = [ ] ;
12+ private finished = false ;
13+ private abortController ?: AbortController ;
14+ private maxQueue : number ;
15+ private overflow : 'ignore' | 'throw' ;
16+
17+ constructor (
18+ private fs : any ,
19+ private path : misc . PathLike ,
20+ private options : opts . IWatchOptions = { } ,
21+ ) {
22+ this . maxQueue = options . maxQueue || 2048 ;
23+ this . overflow = options . overflow || 'ignore' ;
24+ this . startWatching ( ) ;
25+
26+ // Handle AbortSignal
27+ if ( options . signal ) {
28+ if ( options . signal . aborted ) {
29+ this . finish ( ) ;
30+ return ;
31+ }
32+ options . signal . addEventListener ( 'abort' , ( ) => {
33+ this . finish ( ) ;
34+ } ) ;
35+ }
36+ }
37+
38+ private startWatching ( ) {
39+ try {
40+ this . watcher = this . fs . watch ( this . path , this . options , ( eventType : string , filename : string ) => {
41+ this . enqueueEvent ( { eventType, filename } ) ;
42+ } ) ;
43+ } catch ( error ) {
44+ // If we can't start watching, finish immediately
45+ this . finish ( ) ;
46+ throw error ;
47+ }
48+ }
49+
50+ private enqueueEvent ( event : { eventType : string ; filename : string | Buffer } ) {
51+ if ( this . finished ) return ;
52+
53+ // Handle queue overflow
54+ if ( this . eventQueue . length >= this . maxQueue ) {
55+ if ( this . overflow === 'throw' ) {
56+ const error = new Error ( `Watch queue overflow: more than ${ this . maxQueue } events queued` ) ;
57+ this . finish ( error ) ;
58+ return ;
59+ } else {
60+ // 'ignore' - drop the oldest event
61+ this . eventQueue . shift ( ) ;
62+ console . warn ( `Watch queue overflow: dropping event due to exceeding maxQueue of ${ this . maxQueue } ` ) ;
63+ }
64+ }
65+
66+ this . eventQueue . push ( event ) ;
67+
68+ // If there's a waiting promise, resolve it
69+ if ( this . resolveQueue . length > 0 ) {
70+ const { resolve } = this . resolveQueue . shift ( ) ! ;
71+ const nextEvent = this . eventQueue . shift ( ) ! ;
72+ resolve ( { value : nextEvent , done : false } ) ;
73+ }
74+ }
75+
76+ private finish ( error ?: Error ) {
77+ if ( this . finished ) return ;
78+ this . finished = true ;
79+
80+ if ( this . watcher ) {
81+ this . watcher . close ( ) ;
82+ this . watcher = null ;
83+ }
84+
85+ // Resolve or reject all pending promises
86+ while ( this . resolveQueue . length > 0 ) {
87+ const { resolve, reject } = this . resolveQueue . shift ( ) ! ;
88+ if ( error ) {
89+ reject ( error ) ;
90+ } else {
91+ resolve ( { value : undefined , done : true } ) ;
92+ }
93+ }
94+ }
95+
96+ async next ( ) : Promise < IteratorResult < { eventType : string ; filename : string | Buffer } > > {
97+ if ( this . finished ) {
98+ return { value : undefined , done : true } ;
99+ }
100+
101+ // If we have queued events, return one
102+ if ( this . eventQueue . length > 0 ) {
103+ const event = this . eventQueue . shift ( ) ! ;
104+ return { value : event , done : false } ;
105+ }
106+
107+ // Otherwise, wait for the next event
108+ return new Promise ( ( resolve , reject ) => {
109+ this . resolveQueue . push ( { resolve, reject } ) ;
110+ } ) ;
111+ }
112+
113+ async return ( ) : Promise < IteratorResult < { eventType : string ; filename : string | Buffer } > > {
114+ this . finish ( ) ;
115+ return { value : undefined , done : true } ;
116+ }
117+
118+ async throw ( error : any ) : Promise < IteratorResult < { eventType : string ; filename : string | Buffer } > > {
119+ this . finish ( error ) ;
120+ throw error ;
121+ }
122+
123+ [ Symbol . asyncIterator ] ( ) : AsyncIterableIterator < { eventType : string ; filename : string | Buffer } > {
124+ return this ;
125+ }
126+ }
127+
7128export class FsPromises implements FsPromisesApi {
8129 public readonly constants = constants ;
9130
@@ -72,7 +193,11 @@ export class FsPromises implements FsPromisesApi {
72193 ) ;
73194 } ;
74195
75- public readonly watch = ( ) => {
76- throw new Error ( 'Not implemented' ) ;
196+ public readonly watch = (
197+ filename : misc . PathLike ,
198+ options ?: opts . IWatchOptions | string ,
199+ ) : AsyncIterableIterator < { eventType : string ; filename : string | Buffer } > => {
200+ const watchOptions : opts . IWatchOptions = typeof options === 'string' ? { encoding : options as any } : options || { } ;
201+ return new FSWatchAsyncIterator ( this . fs , filename , watchOptions ) ;
77202 } ;
78203}
0 commit comments