Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 9 additions & 6 deletions API.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions docs/REQUIREMENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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() | ✅ |

Expand Down
5 changes: 4 additions & 1 deletion index.bs
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -836,7 +837,9 @@ The <dfn method for="Stream">pipeToSync(source, ...args)</dfn> method is the syn
<li>Let |preventClose| be |options|["{{PipeToSyncOptions/preventClose}}"] and |preventFail| be |options|["{{PipeToSyncOptions/preventFail}}"].
<li>Let |pipeline| be the result of [=compose sync transform pipeline=] with |normalized| and the extracted |transforms|.
<li>Let |totalBytes| be 0.
<li>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.
<li>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.
<li>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|.
<li>For each successfully written |chunk|, add |chunk|'s byte length to |totalBytes|.
<li>On successful completion, if |preventClose| is `false` and |writer| has an `endSync` method, call |writer|.{{SyncWriter/endSync()}}.
<li>On error |e|, if |preventFail| is `false` and |writer| has a `fail` method, call |writer|.{{SyncWriter/fail()}} with |e|. Re-throw |e|.
<li>Return |totalBytes|.
Expand Down
17 changes: 11 additions & 6 deletions src/pull.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
}
}
Expand Down
41 changes: 41 additions & 0 deletions src/pull.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<SyncWriter>).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]() {
Expand Down
15 changes: 10 additions & 5 deletions src/pull.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
}
}
Expand Down
8 changes: 3 additions & 5 deletions src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
*/
Expand All @@ -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(). */
Expand Down
Loading