@@ -22,23 +22,16 @@ async function collect(readable) {
2222 return Buffer . concat ( chunks ) ;
2323}
2424
25- // Disposal of a streaming source completes asynchronously - its descriptor is
26- // closed on the libuv threadpool - and an abandoned archive tears its queued
27- // sources down one after another, each waiting on the previous close. So the
28- // moment "every source is destroyed" can trail the destroy() call, arbitrarily
29- // far on a loaded machine. Wait for that real end state rather than assuming a
30- // fixed delay has been long enough.
31- async function waitForAllDestroyed ( streams ) {
32- const deadline = Date . now ( ) + common . platformTimeout ( 5000 ) ;
33- while ( streams . some ( ( s ) => ! s . destroyed ) && Date . now ( ) < deadline ) {
34- await new Promise ( ( resolve ) => setTimeout ( resolve , 5 ) ) ;
35- }
25+ // destroy() sets .destroyed before a pending open or close finishes. Wait for
26+ // the sources to close before removing their files.
27+ async function waitForAllClosed ( streams ) {
28+ await Promise . all ( streams . map ( ( stream ) =>
29+ ( stream . closed ? undefined : new Promise ( ( resolve ) => stream . once ( 'close' , resolve ) ) ) ) ) ;
3630}
3731
38- // Write `count` throwaway files and return streaming entries backed by their
39- // (eagerly opened) read streams, so a test can observe whether each source is
40- // destroyed. The worst case for leaks: every descriptor is open before the
41- // archive starts.
32+ // Write `count` throwaway files and return entries backed by read streams,
33+ // so the test can observe how disposal closes every source. Their asynchronous
34+ // opens may still be pending when the archive starts.
4235async function streamingEntries ( count ) {
4336 const dir = await fsp . mkdtemp ( path . join ( tmpdir . path , 'zlib-zip-dispose-' ) ) ;
4437 const streams = [ ] ;
@@ -166,7 +159,7 @@ test('an abandoned createZipArchive() destroys the sources of entries it never r
166159 if ( seen > 256 * 1024 ) break ; // Bail while still inside the first member
167160 }
168161 archive . destroy ( ) ;
169- await waitForAllDestroyed ( streams ) ;
162+ await waitForAllClosed ( streams ) ;
170163 const open = streams . filter ( ( s ) => ! s . destroyed ) ;
171164 assert . strictEqual ( open . length , 0 , `${ open . length } source streams left open` ) ;
172165 } finally {
@@ -196,7 +189,7 @@ test('zipEntry disposal destroys a streaming source, sync and async', async () =
196189 try {
197190 entries [ 0 ] [ Symbol . dispose ] ( ) ;
198191 await entries [ 1 ] [ Symbol . asyncDispose ] ( ) ;
199- await waitForAllDestroyed ( streams ) ;
192+ await waitForAllClosed ( streams ) ;
200193 assert . strictEqual ( streams [ 0 ] . destroyed , true ) ;
201194 assert . strictEqual ( streams [ 1 ] . destroyed , true ) ;
202195 } finally {
@@ -232,7 +225,7 @@ test('createZipArchiveSync() throws on a streaming entry and disposes the rest',
232225 try {
233226 assert . throws ( ( ) => Array . from ( zlib . createZipArchiveSync ( entries ) ) ,
234227 { code : 'ERR_INVALID_STATE' } ) ;
235- await waitForAllDestroyed ( streams ) ;
228+ await waitForAllClosed ( streams ) ;
236229 assert . strictEqual ( streams . filter ( ( s ) => ! s . destroyed ) . length , 0 ) ;
237230 } finally {
238231 await fsp . rm ( dir , { recursive : true , force : true } ) ;
0 commit comments