From 97b03dda0b8ffb96c2d0256cddbdf34a86f8b0a5 Mon Sep 17 00:00:00 2001 From: "Kamat, Trivikram" <16024985+trivikr@users.noreply.github.com> Date: Sun, 31 May 2026 14:15:13 -0700 Subject: [PATCH] fix: fail pipeToSync when sync writer returns false --- API.md | 15 +++++++++------ docs/REQUIREMENTS.md | 2 ++ index.bs | 5 ++++- src/pull.js | 17 +++++++++++------ src/pull.test.ts | 41 +++++++++++++++++++++++++++++++++++++++++ src/pull.ts | 15 ++++++++++----- src/types.ts | 8 +++----- 7 files changed, 80 insertions(+), 23 deletions(-) diff --git a/API.md b/API.md index f5052b7..0fc03e1 100644 --- a/API.md +++ b/API.md @@ -286,12 +286,10 @@ interface SyncWriter { } ``` -`SyncWriter` follows the same backpressure policies as `Writer`: -- `"block"`: `write()`/`writev()` enqueue and return `true` when buffer has - space, or enqueue and return `false` when full (backpressure signal; data IS - accepted). -- `"strict"`: Writes exceeding buffer capacity throw a `RangeError`. -- `"drop-oldest"`/`"drop-newest"`: Behave as with `Writer`. Writes never fail. +For `SyncWriter`, `write()` and `writev()` return `true` only when the data was +accepted synchronously. A `false` return means the operation could not be +performed synchronously, so sync callers such as `Stream.pipeToSync()` must +treat it as a failed write. Failed writes are not counted as written bytes. - `end()` throws `TypeError` if already closed or errored (no -1 fallback; there is no async counterpart to fall back to). - `fail()` transitions to error state; no-op if already closed or errored. @@ -735,6 +733,11 @@ function pipeToSync( ): number ``` +`pipeToSync()` counts bytes only after a successful synchronous write. If +`writer.write()` or `writer.writev()` returns `false`, `pipeToSync()` throws, +calls `writer.fail()` unless `preventFail` is set, and does not count that +failed write. + --- ## Consumers diff --git a/docs/REQUIREMENTS.md b/docs/REQUIREMENTS.md index 48064b4..150b49f 100644 --- a/docs/REQUIREMENTS.md +++ b/docs/REQUIREMENTS.md @@ -263,6 +263,8 @@ Consumes source and writes to a writer with optional transforms. |----|-------------|--------| | WRITE-020 | Fails writer on source error | ✅ | | WRITE-021 | Throws if no writer provided | ✅ | +| WRITE-022 | pipeToSync() throws and fails writer when writev returns false | ✅ | +| WRITE-023 | pipeToSync() throws and fails writer when write returns false | ✅ | | WRITE-025 | pipeTo passes signal to writer.write() | ✅ | | WRITE-026 | pipeTo passes signal to writer.end() | ✅ | diff --git a/index.bs b/index.bs index df2c4ef..16b9712 100644 --- a/index.bs +++ b/index.bs @@ -375,6 +375,7 @@ Like {{Writer}}, {{SyncWriter}} is an interface, not a concrete class. Implement * With `"strict"` and `"unbounded"`, {{SyncWriter/writeSync()}} and {{SyncWriter/writevSync()}} enqueue the chunk(s) and return `true` when the [=byte budget=] has capacity, or return `false` when the budget is exhausted. The data is not accepted; no async fallback is available. * With `"drop-oldest"` and `"drop-newest"`, writes behave as described in [[#concept-backpressure]]: the oldest or newest data is discarded respectively. Writes never fail. +* {{SyncWriter/writeSync()}} and {{SyncWriter/writevSync()}} return `true` only when the data was accepted synchronously. A `false` return means the operation could not be performed synchronously; synchronous callers such as {{Stream/pipeToSync()}} must treat it as a failed write. * {{SyncWriter/endSync()}} throws a {{TypeError}} if the writer is already closed or errored. * {{SyncWriter/fail()}} transitions the writer to the errored state. If the writer is already closed or errored, it is a no-op. @@ -836,7 +837,9 @@ The pipeToSync(source, ...args) method is the syn
  • Let |preventClose| be |options|["{{PipeToSyncOptions/preventClose}}"] and |preventFail| be |options|["{{PipeToSyncOptions/preventFail}}"].
  • Let |pipeline| be the result of [=compose sync transform pipeline=] with |normalized| and the extracted |transforms|.
  • Let |totalBytes| be 0. -
  • Synchronously iterate |pipeline|. For each |batch| yielded, for each |chunk| in |batch|, add |chunk|'s byte length to |totalBytes|. Then write the batch to |writer|: if |writer| has a `writevSync` method, call |writer|.{{SyncWriter/writevSync()}} with |batch|; otherwise call |writer|.{{SyncWriter/writeSync()}} for each chunk. If the write returns `false`, throw a {{RangeError}} indicating the destination's [=byte budget=] is exhausted. +
  • Synchronously iterate |pipeline|. For each |batch| yielded, write the batch to |writer|: if |writer| has a `writevSync` method, call |writer|.{{SyncWriter/writevSync()}} with |batch|; otherwise call |writer|.{{SyncWriter/writeSync()}} for each chunk. +
  • If any call to {{SyncWriter/writeSync()}} or {{SyncWriter/writevSync()}} returns `false`, throw a {{RangeError}} indicating the destination's [=byte budget=] is exhausted, and do not add the failed write's byte length to |totalBytes|. +
  • For each successfully written |chunk|, add |chunk|'s byte length to |totalBytes|.
  • On successful completion, if |preventClose| is `false` and |writer| has an `endSync` method, call |writer|.{{SyncWriter/endSync()}}.
  • On error |e|, if |preventFail| is `false` and |writer| has a `fail` method, call |writer|.{{SyncWriter/fail()}} with |e|. Re-throw |e|.
  • Return |totalBytes|. diff --git a/src/pull.js b/src/pull.js index 2ab1857..123b199 100644 --- a/src/pull.js +++ b/src/pull.js @@ -1206,17 +1206,22 @@ function pipeToSync(source) { try { for (var _c = 0, pipeline_1 = pipeline; _c < pipeline_1.length; _c++) { var batch = pipeline_1[_c]; - for (var _d = 0, batch_1 = batch; _d < batch_1.length; _d++) { - var chunk = batch_1[_d]; - totalBytes += chunk.byteLength; - } if ('writev' in writer && typeof writer.writev === 'function') { - writer.writev(batch); + if (!writer.writev(batch)) { + throw new Error('Sync writev failed'); + } + for (var _d = 0, batch_1 = batch; _d < batch_1.length; _d++) { + var chunk = batch_1[_d]; + totalBytes += chunk.byteLength; + } } else { for (var _e = 0, batch_2 = batch; _e < batch_2.length; _e++) { var chunk = batch_2[_e]; - writer.write(chunk); + if (!writer.write(chunk)) { + throw new Error('Sync write failed'); + } + totalBytes += chunk.byteLength; } } } diff --git a/src/pull.test.ts b/src/pull.test.ts index 27d3cc8..c16f6a8 100644 --- a/src/pull.test.ts +++ b/src/pull.test.ts @@ -47,11 +47,13 @@ function createMockSyncWriter(): SyncWriter & { chunks: Uint8Array[]; closed: bo if (this.closed) throw new Error('Writer is closed'); const data = typeof chunk === 'string' ? new TextEncoder().encode(chunk) : chunk; this.chunks.push(data); + return true; }, writev(chunks: (Uint8Array | string)[]) { for (const chunk of chunks) { this.write(chunk); } + return true; }, end() { this.closed = true; @@ -762,6 +764,45 @@ describe('pipeToSync()', () => { }); describe('error handling', () => { + it('should throw and fail writer when writev returns false [WRITE-022]', () => { + const source = fromSync('ab'); + const writer = { + ...createMockSyncWriter(), + writev() { + return false; + }, + }; + + assert.throws(() => { + pipeToSync(source, writer); + }, /Sync writev failed/); + + assert.ok(writer.failed); + assert.deepStrictEqual(writer.chunks, []); + }); + + it('should throw and fail writer when write returns false [WRITE-023]', () => { + const source = fromSync(['a', 'b']); + const writer = createMockSyncWriter(); + delete (writer as Partial).writev; + writer.write = function write(chunk: Uint8Array | string) { + if (this.chunks.length >= 1) { + return false; + } + const data = typeof chunk === 'string' ? new TextEncoder().encode(chunk) : chunk; + this.chunks.push(data); + return true; + }; + + assert.throws(() => { + pipeToSync(source, writer); + }, /Sync write failed/); + + assert.ok(writer.failed); + assert.strictEqual(writer.chunks.length, 1); + assert.strictEqual(decode(concatBytes(writer.chunks)), 'a'); + }); + it('should fail writer on source error [WRITE-020]', () => { const source = { *[Symbol.iterator]() { diff --git a/src/pull.ts b/src/pull.ts index 6fa2aab..4b3b236 100644 --- a/src/pull.ts +++ b/src/pull.ts @@ -712,14 +712,19 @@ export function pipeToSync( try { for (const batch of pipeline) { - for (const chunk of batch) { - totalBytes += chunk.byteLength; - } if ('writev' in writer && typeof writer.writev === 'function') { - writer.writev(batch); + if (!writer.writev(batch)) { + throw new Error('Sync writev failed'); + } + for (const chunk of batch) { + totalBytes += chunk.byteLength; + } } else { for (const chunk of batch) { - writer.write(chunk); + if (!writer.write(chunk)) { + throw new Error('Sync write failed'); + } + totalBytes += chunk.byteLength; } } } diff --git a/src/types.ts b/src/types.ts index f6449aa..9cb764d 100644 --- a/src/types.ts +++ b/src/types.ts @@ -203,10 +203,8 @@ export interface Writer { /** * Sync writer interface for producing data. * - * Follows the same backpressure policies as Writer: - * - "block": write/writev enqueue and return true (space) or false (backpressure signal; data IS accepted) - * - "strict": writes exceeding buffer capacity throw RangeError - * - "drop-oldest"/"drop-newest": writes never fail + * A false return from write/writev means the writer could not perform the + * operation synchronously. Sync pipelines must treat that as a failed write. * - end() throws TypeError if already closed or errored * - fail() is a no-op if already closed or errored */ @@ -219,7 +217,7 @@ export interface SyncWriter { */ readonly desiredSize: number | null; - /** Write single chunk. Returns true if buffer has space, false as backpressure signal under "block". */ + /** Write single chunk. Returns true if accepted, false if not completed synchronously. */ write(chunk: Uint8Array | string): boolean; /** Write multiple chunks atomically. All-or-nothing. Returns true/false like write(). */