| 507 | // Wrap `stream` so its capnweb envelope is released exactly once on clean |
| 508 | // completion, source failure, or consumer cancellation. |
| 509 | function disposeOnDone<T>(stream: ReadableStream<T>, onDone: () => void): ReadableStream<T> { |
| 510 | const reader = stream.getReader(); |
| 511 | let finished = false; |
| 512 | const finish = () => { |
| 513 | if (finished) return; |
| 514 | finished = true; |
| 515 | try { |
| 516 | reader.releaseLock(); |
| 517 | } catch {} |
| 518 | try { |
| 519 | onDone(); |
| 520 | } catch { |
| 521 | // Disposer failures cannot be recovered at this boundary. |
| 522 | } |
| 523 | }; |
| 524 | return new ReadableStream<T>( |
| 525 | { |
| 526 | async pull(controller) { |
| 527 | try { |
| 528 | const { value, done } = await reader.read(); |
| 529 | if (done) { |
| 530 | finish(); |
| 531 | controller.close(); |
| 532 | } else { |
| 533 | controller.enqueue(value); |
| 534 | } |
| 535 | } catch (error) { |
| 536 | finish(); |
| 537 | controller.error(error); |
| 538 | } |
| 539 | }, |
| 540 | async cancel(reason) { |
| 541 | try { |
| 542 | await reader.cancel(reason); |
| 543 | } finally { |
| 544 | finish(); |
| 545 | } |
| 546 | }, |
| 547 | }, |
| 548 | { highWaterMark: 0 }, |
| 549 | ); |
| 550 | } |
| 551 | |
| 552 | // Best-effort dispose of a capnweb result envelope. Real envelopes |
| 553 | // expose [Symbol.dispose]; test fakes return plain objects, so the |