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(). */