11'use strict' ;
22// Flags: --no-warnings --expose-internals
3- require ( '../common' ) ;
3+ const common = require ( '../common' ) ;
44const assert = require ( 'assert' ) ;
55const test = require ( 'node:test' ) ;
66const { Duplex, Writable } = require ( 'stream' ) ;
99 newReadableWritablePairFromDuplex,
1010} = require ( 'internal/webstreams/adapters' ) ;
1111
12+ function isSameError ( expected ) {
13+ return common . mustCall ( ( actual ) => {
14+ assert . strictEqual ( actual , expected ) ;
15+ return true ;
16+ } ) ;
17+ }
18+
1219// Verify that when the underlying Node.js stream throws synchronously from
1320// write(), the writable web stream properly rejects but does not destroy
1421// the stream (destroy-on-sync-throw is only used internally by
@@ -34,6 +41,72 @@ test('WritableStream from Node.js stream handles sync write throw', async () =>
3441 assert . strictEqual ( writable . destroyed , false ) ;
3542} ) ;
3643
44+ test ( 'WritableStream from Node.js stream handles async write error' , async ( ) => {
45+ const error = new Error ( 'boom' ) ;
46+ const writable = new Writable ( {
47+ write ( _chunk , _encoding , callback ) {
48+ setImmediate ( callback , error ) ;
49+ } ,
50+ } ) ;
51+ const writer = Writable . toWeb ( writable ) . getWriter ( ) ;
52+
53+ await Promise . all ( [
54+ assert . rejects ( writer . write ( Buffer . from ( 'hello' ) ) , isSameError ( error ) ) ,
55+ assert . rejects ( writer . closed , isSameError ( error ) ) ,
56+ ] ) ;
57+ } ) ;
58+
59+ test ( 'WritableStream aborts while a native write is pending' , async ( ) => {
60+ const error = new Error ( 'abort' ) ;
61+ let finishWrite ;
62+ let startWrite ;
63+ const writeStarted = new Promise ( ( resolve ) => {
64+ startWrite = resolve ;
65+ } ) ;
66+ const writable = new Writable ( {
67+ write ( _chunk , _encoding , callback ) {
68+ finishWrite = callback ;
69+ startWrite ( ) ;
70+ } ,
71+ } ) ;
72+ const writer = Writable . toWeb ( writable ) . getWriter ( ) ;
73+ const writePromise = writer . write ( Buffer . from ( 'hello' ) ) ;
74+ await writeStarted ;
75+
76+ const writeRejected = assert . rejects ( writePromise , isSameError ( error ) ) ;
77+ const closedRejected = assert . rejects ( writer . closed , isSameError ( error ) ) ;
78+ await Promise . all ( [
79+ writer . abort ( error ) ,
80+ writeRejected ,
81+ closedRejected ,
82+ ] ) ;
83+
84+ finishWrite ( ) ;
85+ await new Promise ( setImmediate ) ;
86+ assert . strictEqual ( writable . destroyed , true ) ;
87+ } ) ;
88+
89+ test ( 'WritableStream handles destruction while a write is pending' , async ( ) => {
90+ const error = new Error ( 'destroy' ) ;
91+ let startWrite ;
92+ const writeStarted = new Promise ( ( resolve ) => {
93+ startWrite = resolve ;
94+ } ) ;
95+ const writable = new Writable ( {
96+ write ( _chunk , _encoding , _callback ) {
97+ startWrite ( ) ;
98+ } ,
99+ } ) ;
100+ const writer = Writable . toWeb ( writable ) . getWriter ( ) ;
101+ const writePromise = writer . write ( Buffer . from ( 'hello' ) ) ;
102+ await writeStarted ;
103+
104+ const writeRejected = assert . rejects ( writePromise , isSameError ( error ) ) ;
105+ const closedRejected = assert . rejects ( writer . closed , isSameError ( error ) ) ;
106+ writable . destroy ( error ) ;
107+ await Promise . all ( [ writeRejected , closedRejected ] ) ;
108+ } ) ;
109+
37110test ( 'Duplex-backed pair does NOT destroy on sync write throw' , async ( ) => {
38111 const error = new TypeError ( 'invalid chunk' ) ;
39112 const duplex = new Duplex ( {
0 commit comments