@@ -17,13 +17,11 @@ const {
1717 ObjectDefineProperties,
1818 ObjectKeys,
1919 PromisePrototypeThen,
20- PromiseResolve,
2120 PromiseWithResolvers,
2221 SafeSet,
2322 Symbol,
2423 SymbolAsyncDispose,
2524 SymbolAsyncIterator,
26- SymbolDispose,
2725 SymbolIterator,
2826 TypedArrayPrototypeGetByteLength,
2927 Uint8Array,
@@ -134,16 +132,9 @@ const {
134132} = require ( 'internal/blob' ) ;
135133
136134const {
137- drainableProtocol,
138135 kValidatedSource,
139136} = require ( 'internal/streams/iter/types' ) ;
140137
141- const {
142- convertChunks,
143- getWriterSignal,
144- toWriterUint8Array,
145- } = require ( 'internal/streams/iter/utils' ) ;
146-
147138const {
148139 from : streamFrom ,
149140 fromSync : streamFromSync ,
@@ -232,6 +223,10 @@ const {
232223 QuicStreamState,
233224} = require ( 'internal/quic/state' ) ;
234225
226+ const {
227+ getWriter,
228+ } = require ( 'internal/quic/writer' ) ;
229+
235230const assert = require ( 'internal/assert' ) ;
236231
237232const {
@@ -1559,328 +1554,6 @@ function isSyncIterable(obj) {
15591554 return obj != null && typeof obj [ SymbolIterator ] === 'function' ;
15601555}
15611556
1562- function getWriter ( obj ) {
1563- const {
1564- setDrainCallback,
1565- setStopSendingCallback,
1566- writeDesiredSize,
1567- writeEnded,
1568- writeChunks,
1569- resetStream,
1570- endWrite,
1571- nonWritableStream,
1572- initStreamingSource,
1573- } = obj ;
1574- // TODO, if this is not an internal function,
1575- // the input should be checked for types
1576- let closed = false ;
1577- let ending = false ;
1578- let errored = false ;
1579- let error = null ;
1580- let totalBytesWritten = 0 ;
1581- let drainWakeup = null ;
1582- let pendingWrite = null ;
1583- let pendingEnd = null ;
1584-
1585- function raceWithSignal ( promise , signal ) {
1586- if ( signal === undefined ) return promise ;
1587- signal . throwIfAborted ( ) ;
1588-
1589- const deferred = PromiseWithResolvers ( ) ;
1590- const onAbort = ( ) => deferred . reject ( signal . reason ) ;
1591- signal . addEventListener ( 'abort' , onAbort , {
1592- __proto__ : null ,
1593- once : true ,
1594- } ) ;
1595- PromisePrototypeThen (
1596- promise ,
1597- ( value ) => {
1598- signal . removeEventListener ( 'abort' , onAbort ) ;
1599- deferred . resolve ( value ) ;
1600- } ,
1601- ( reason ) => {
1602- signal . removeEventListener ( 'abort' , onAbort ) ;
1603- deferred . reject ( reason ) ;
1604- } ) ;
1605- return deferred . promise ;
1606- }
1607-
1608- function waitForDrain ( signal ) {
1609- drainWakeup ??= PromiseWithResolvers ( ) ;
1610- return raceWithSignal ( drainWakeup . promise , signal ) ;
1611- }
1612-
1613- // Drain callback - The implementation C++/js fires this when send buffer has space
1614- setDrainCallback ( ( ) => {
1615- if ( drainWakeup ) {
1616- drainWakeup . resolve ( true ) ;
1617- drainWakeup = null ;
1618- }
1619- } ) ;
1620-
1621- setStopSendingCallback ( ( reason ) => {
1622- if ( ! closed && ! errored ) {
1623- errored = true ;
1624- error = reason ;
1625- if ( drainWakeup != null ) {
1626- markPromiseAsHandled ( drainWakeup . promise ) ;
1627- drainWakeup . reject ( error ) ;
1628- drainWakeup = null ;
1629- }
1630- }
1631- } ) ;
1632-
1633- // A note on backpressure handling: per the stream/iter spec, the default
1634- // backpressure policy for writers is strict. One async write may wait for
1635- // capacity; additional writes are rejected until it settles.
1636-
1637- function writeConvertedSync ( chunk , token ) {
1638- const isPending = pendingWrite === token && token !== undefined ;
1639- // If the stream is closed, errored, or write-ended, we cannot accept
1640- // more data. Refuse the sync write.
1641- if ( closed || errored || writeEnded ( ) ||
1642- ( ending && ! isPending ) ||
1643- ( pendingWrite !== null && ! isPending ) ) {
1644- return false ;
1645- }
1646- const len = TypedArrayPrototypeGetByteLength ( chunk ) ;
1647- if ( len === 0 ) return true ;
1648- // Refuse the write only when there is no available capacity at
1649- // all. If we can write we allow the write even if the
1650- // chunk is larger than the remaining capacity -- the source
1651- // will accept the data into the underlying queues e.g. DataQueue and
1652- // UpdateWriteDesiredSize() will drop writeDesiredSize toward 0,
1653- // at which point the standard drain mechanism takes over.
1654- // This follows the iter-streams model where writes beyond the
1655- // budget succeed and backpressure applies to *subsequent* writes.
1656- if ( writeDesiredSize ( ) === 0 ) return false ;
1657- const result = writeChunks ( [ chunk ] ) ;
1658- if ( result === undefined ) return false ;
1659- totalBytesWritten += len ;
1660- return true ;
1661- }
1662-
1663- function writeSync ( chunk ) {
1664- return writeConvertedSync ( toWriterUint8Array ( chunk ) ) ;
1665- }
1666-
1667- function write ( chunk , options = kEmptyObject ) {
1668- chunk = toWriterUint8Array ( chunk ) ;
1669- const signal = getWriterSignal ( options ) ;
1670- return writeAsync ( chunk , signal ) ;
1671- }
1672-
1673- async function writeAsync ( chunk , signal ) {
1674- if ( errored ) throw error ;
1675- if ( closed || ending || writeEnded ( ) ) {
1676- throw new ERR_INVALID_STATE ( 'Writer is closed' ) ;
1677- }
1678- signal ?. throwIfAborted ( ) ;
1679-
1680- if ( writeConvertedSync ( chunk ) ) return ;
1681- await writeWhenDrained ( chunk , writeConvertedSync , signal ) ;
1682- }
1683-
1684- async function writeWhenDrained ( chunks , writeConverted , signal ) {
1685- if ( pendingWrite !== null ) {
1686- throw new ERR_INVALID_STATE . RangeError (
1687- 'Backpressure violation: too many pending writes. ' +
1688- 'Await each write() call to respect backpressure.' ) ;
1689- }
1690- if ( writeDesiredSize ( ) !== 0 ) {
1691- throw new ERR_INVALID_STATE ( 'Stream write buffer is full' ) ;
1692- }
1693-
1694- const done = PromiseWithResolvers ( ) ;
1695- const token = { __proto__ : null , done } ;
1696- pendingWrite = token ;
1697- try {
1698- while ( true ) {
1699- await waitForDrain ( signal ) ;
1700- if ( errored ) throw error ;
1701- if ( closed || writeEnded ( ) ) {
1702- throw new ERR_INVALID_STATE ( 'Writer is closed' ) ;
1703- }
1704- signal ?. throwIfAborted ( ) ;
1705- if ( writeConverted ( chunks , token ) ) return ;
1706- if ( writeDesiredSize ( ) !== 0 ) {
1707- throw new ERR_INVALID_STATE ( 'Stream write buffer is full' ) ;
1708- }
1709- }
1710- } finally {
1711- if ( pendingWrite === token ) pendingWrite = null ;
1712- done . resolve ( ) ;
1713- }
1714- }
1715-
1716- function writevConvertedSync ( chunks , token ) {
1717- const isPending = pendingWrite === token && token !== undefined ;
1718- if ( closed || errored || writeEnded ( ) ||
1719- ( ending && ! isPending ) ||
1720- ( pendingWrite !== null && ! isPending ) ) {
1721- return false ;
1722- }
1723- let len = 0 ;
1724- for ( const c of chunks ) len += TypedArrayPrototypeGetByteLength ( c ) ;
1725- if ( len === 0 ) return true ;
1726- if ( writeDesiredSize ( ) === 0 ) return false ;
1727- const result = writeChunks ( chunks ) ;
1728- if ( result === undefined ) return false ;
1729- totalBytesWritten += len ;
1730- return true ;
1731- }
1732-
1733- function writevSync ( chunks ) {
1734- return writevConvertedSync ( convertChunks ( chunks ) ) ;
1735- }
1736-
1737- function writev ( chunks , options = kEmptyObject ) {
1738- chunks = convertChunks ( chunks ) ;
1739- const signal = getWriterSignal ( options ) ;
1740- return writevAsync ( chunks , signal ) ;
1741- }
1742-
1743- async function writevAsync ( chunks , signal ) {
1744- if ( errored ) throw error ;
1745- if ( closed || ending || writeEnded ( ) ) {
1746- throw new ERR_INVALID_STATE ( 'Writer is closed' ) ;
1747- }
1748- signal ?. throwIfAborted ( ) ;
1749-
1750- if ( writevConvertedSync ( chunks ) ) return ;
1751- await writeWhenDrained ( chunks , writevConvertedSync , signal ) ;
1752- }
1753-
1754- function endSync ( ) {
1755- // Per the streams/iter spec, endSync and end follow a try-fallback
1756- // pattern. That is, callers should try endSync first and if it returns
1757- // -1, then they should call and await end(). This is a signal that sync
1758- // end is not currently possible. However, we always support sync end
1759- // here unless the stream is already errored.
1760- if ( errored ) return - 1 ;
1761-
1762- // If we're already closed, just return the total bytes written.
1763- if ( closed ) return totalBytesWritten ;
1764-
1765- // Accepted writes and existing drain waiters must settle before ending.
1766- if ( ending || pendingWrite !== null || drainWakeup !== null ) return - 1 ;
1767-
1768- // Fantastic, we can end synchronously!
1769- endWrite ( ) ;
1770- closed = true ;
1771- return totalBytesWritten ;
1772- }
1773-
1774- function end ( options = kEmptyObject ) {
1775- const signal = getWriterSignal ( options ) ;
1776- return endAsync ( signal ) ;
1777- }
1778-
1779- async function endAsync ( signal ) {
1780- if ( errored ) throw error ;
1781- if ( closed ) return totalBytesWritten ;
1782- signal ?. throwIfAborted ( ) ;
1783-
1784- // Per the streams/iter spec, endSync and end follow a try-fallback
1785- // pattern. That is, callers should try endSync first and if it returns
1786- // -1, then they should call and await end(). This is a signal that sync
1787- // end is not currently possible. However, we always support sync end
1788- // here unless the stream is already errored.
1789- // While the user should have already called endSync, we call it again
1790- // here to actually process the end request. At worst it's called twice.
1791- const n = endSync ( ) ;
1792-
1793- // A return value of -1 indicates that endSync was not yet able to
1794- // process the end request, either because we are errored or because we
1795- // are awaiting drain. If we're errored, throw the error. If we're waiting
1796- // for drain, await it and then try ending again.
1797-
1798- if ( n >= 0 ) return n ;
1799- if ( errored ) throw error ;
1800-
1801- if ( pendingEnd === null ) {
1802- ending = true ;
1803- pendingEnd = finishEnd ( ) ;
1804- }
1805- return raceWithSignal ( pendingEnd , signal ) ;
1806- }
1807-
1808- async function finishEnd ( ) {
1809- if ( pendingWrite !== null ) {
1810- await pendingWrite . done . promise ;
1811- }
1812- if ( errored ) throw error ;
1813-
1814- if ( drainWakeup !== null ) {
1815- await drainWakeup . promise ;
1816- }
1817- if ( errored ) throw error ;
1818-
1819- if ( ! writeEnded ( ) ) endWrite ( ) ;
1820- closed = true ;
1821- ending = false ;
1822- return totalBytesWritten ;
1823- }
1824-
1825- function fail ( reason ) {
1826- if ( closed || errored ) return ;
1827- errored = true ;
1828- error = reason ;
1829- resetStream ( error ) ;
1830- if ( drainWakeup != null ) {
1831- markPromiseAsHandled ( drainWakeup . promise ) ;
1832- drainWakeup . reject ( error ) ;
1833- drainWakeup = null ;
1834- }
1835- }
1836-
1837- const writer = {
1838- __proto__ : null ,
1839- get canWrite ( ) {
1840- if ( closed || ending || errored || writeEnded ( ) ) {
1841- return null ;
1842- }
1843- return writeDesiredSize ( ) > 0 ;
1844- } ,
1845- writeSync,
1846- write,
1847- writevSync,
1848- writev,
1849- endSync,
1850- end,
1851- fail,
1852- [ drainableProtocol ] ( ) {
1853- if ( closed || ending || errored ) return null ;
1854- // If a drain is already pending, return the existing promise.
1855- if ( drainWakeup != null ) return drainWakeup . promise ;
1856- if ( writeDesiredSize ( ) > 0 ) return null ;
1857- drainWakeup = PromiseWithResolvers ( ) ;
1858- return drainWakeup . promise ;
1859- } ,
1860- [ SymbolAsyncDispose ] ( ) {
1861- if ( ending ) return pendingEnd ;
1862- if ( ! closed && ! errored ) fail ( ) ;
1863- return PromiseResolve ( ) ;
1864- } ,
1865- [ SymbolDispose ] ( ) {
1866- if ( ! closed && ! errored ) fail ( ) ;
1867- } ,
1868- } ;
1869-
1870- // Non-writable stream - return a pre-closed writer.
1871- // A remote unidirectional stream is read-only and has no writable
1872- // side. isLocal distinguishes locally-initiated (writable) from
1873- // remotely-initiated (read-only) uni streams.
1874- if ( nonWritableStream ( ) || writeEnded ( ) ) {
1875- closed = true ;
1876- return writer ;
1877- }
1878-
1879- // Initialize the outbound DataQueue for streaming writes
1880- initStreamingSource ( ) ;
1881- return writer ;
1882- }
1883-
18841557// Functions used specifically for internal or assertion purposes only.
18851558let getQuicStreamState ;
18861559let getQuicSessionState ;
0 commit comments