Skip to content

Bunch of streams fixes - #4251

Merged
Jarred-Sumner merged 36 commits into
mainfrom
jarred/wip
Aug 23, 2023
Merged

Jarred-Sumner merged 36 commits into
mainfrom
jarred/wip

Conversation

@Jarred-Sumner

@Jarred-Sumner Jarred-Sumner commented Aug 22, 2023 •

Copy link
Copy Markdown
Collaborator

What does this PR do?

How did you verify your code works?

There are tests

@github-actions

github-actions Bot commented Aug 22, 2023 •

Copy link
Copy Markdown
Contributor

✅ test failures on linux-x64-baseline have been resolved.

#daf7e4d5049c5caf47dcc40be40ded76ff305502

@github-actions

github-actions Bot commented Aug 22, 2023 •

Copy link
Copy Markdown
Contributor

❌ @cirospaciari 1 files with test failures on linux-x64:

  • test/bundler/esbuild/splitting.test.ts

View test output

#daf7e4d5049c5caf47dcc40be40ded76ff305502

@github-actions

github-actions Bot commented Aug 22, 2023 •

Copy link
Copy Markdown
Contributor

❌ @cirospaciari 4 files with test failures on bun-darwin-aarch64:

  • test/js/bun/test/test-test.test.ts
  • test/js/node/fs/fs.test.ts
  • test/js/node/watch/fs.watch.test.ts
  • test/js/third_party/resvg/bbox.test.js

View test output

#daf7e4d5049c5caf47dcc40be40ded76ff305502

@github-actions

github-actions Bot commented Aug 22, 2023 •

Copy link
Copy Markdown
Contributor

❌ @cirospaciari 8 files with test failures on bun-darwin-x64-baseline:

  • test/bundler/bundler_compile.test.ts
  • test/js/bun/spawn/spawn-streaming-stdin.test.ts
  • test/js/bun/sqlite/sqlite.test.js
  • test/js/node/fs/fs.test.ts
  • test/js/third_party/resvg/bbox.test.js
  • test/js/third_party/webpack/webpack.test.ts
  • test/js/web/timers/setInterval.test.js
  • test/js/web/timers/setTimeout.test.js

View test output

#daf7e4d5049c5caf47dcc40be40ded76ff305502

@Jarred-Sumner Jarred-Sumner changed the title don't do async hooks stuff when unnecessary Bunch of streams fixes Aug 23, 2023
@Jarred-Sumner
Jarred-Sumner merged commit c603857 into main Aug 23, 2023
@Jarred-Sumner
Jarred-Sumner deleted the jarred/wip branch August 23, 2023 21:05
Jarred-Sumner added a commit that referenced this pull request Sep 9, 2026
…se once, call cancel() only on a client abort (#41894)

### What does this PR do?

Fixes four bugs in `type: "direct"` `ReadableStream`s and gives every
whole-body consumer the same rule for when such a stream ends.

Fixes #17175
Closes #41917

#### Bugs fixed

**1. `controller.close()` sent `Content-Length: 0`.**
A `pull()` that wrote and called `close()` in the same tick answered
`200` with an empty body. `HTTPServerWritable::end()` left the buffered
tail to the auto-flusher, but `do_render_stream` frees the sink as soon
as a sync `pull()` returns.
Fix: `end()` (what `controller.close()` reaches) now takes the
`end_from_js()` path. It flushes and ends the response, or parks
`pending_flush` under backpressure. `close(error)` goes to `fail()`,
which drops what is buffered and finishes the sink without ending the
response; the owner then closes it as incomplete.

**2. The chunked last-chunk was written twice.**
When `finalize()` sent a parked tail, `res.end(tail)` wrote `0\r\n\r\n`
without recording `ended_response`, so `end_stream()` wrote it again.
The next keep-alive response then fails to parse
(`HPE_INVALID_CONSTANT`).
Fix: `ended_response` is set in one place, `mark_response_ended()`, at
the point where uWS `end()`/`try_end()` completes. The other paths
assert it.

**3. `cancel(undefined)` ran after every completed response (#17175).**
The sink's `onClose` callback fires both when the source closes itself
and when the peer goes away, and the handler could not tell them apart.
Fix: `onClose` carries `sinkClosed`, true only from
`JSSinkController__onClose`. `directStreamOnClose` calls the source's
`cancel(reason)` only then. `HTTPServerWritable::abort` passes
`CommonAbortReason::ConnectionClosed` as the reason.

**4. An `async pull()` that resolved without `close()` ended differently
per consumer.**
`Bun.serve` finished the response. `Bun.write(file, new
Response(stream))` and `Bun.spawn({ stdin })` closed the fd and dropped
the unflushed tail. `.text()` / `.bytes()` / `Bun.readableStreamTo*()`
never settled. `Bun.write` also freed the `FileSink` with its controller
still attached, a use-after-free at the next GC.
Fix: one rule, below. `FileSink` ends with the flushing `end()`, and
`Bun.write` pipes JS streams through `FileSink::pipe_stream`, which
detaches the controller before the sink is freed.

#### The rule

A consumer that takes the whole body (`Bun.serve`, `fetch` request
bodies, `Bun.write`, `Bun.spawn` stdin, `.text()` / `.bytes()` /
`Bun.readableStreamTo*()`) calls `pull()` once.

- If it returns a promise, the stream stays open while the promise is
pending, ends when it resolves, and errors when it rejects.
- If it returns synchronously without closing, the stream stays open
until `close()` / `end()`.
- `controller.close(error)` fails the stream. Falsy arguments
(`undefined`, `0`, `""`) are a clean close. `end(x)` is always a clean
close.

`readDirectStream` (native sinks) and the one-shot consumers (`.text()`
etc., via `m_closeOnPullSettled`) follow it.

The JS reader path (`getReader()`, `for await`, `pipeTo()`, `tee()`) is
unchanged: `pull()` is a demand signal there, called again for a later
read once the previous call has settled (#33782). A reader has a demand
signal, a whole-body consumer does not.

#### Behavior changes

| Change | Before |
| --- | --- |
| An `async pull()` that resolves without `close()` ends the stream on
every whole-body consumer | `Bun.serve` already did. `.text()` etc.
hung; `Bun.write` truncated or crashed |
| `controller.close()` flushes buffered bytes on the HTTP sink | Sync
`pull()` + `close()` sent `Content-Length: 0` |
| `controller.close(error)` fails the stream everywhere: readers reject,
`Bun.serve` aborts the response with no chunked terminator, `fetch()`
with such a request body rejects | The reader path treated it as a clean
close |
| `cancel(reason)` runs only when the consumer goes away. On a client
disconnect the reason is the connection-closed `AbortError` | Also ran
after the source's own `close()` / `end()`, with `undefined` |
| On a native sink, `controller.write()` / `flush()` return `0` after
the client, file, or child process went away | Threw `"This
HTTPResponseSink has already been closed…"` |
| A `pull()` that rejects after the stream already closed cleanly, or
after the peer left, is ignored | `.bytes()` and `fetch()` request
bodies surfaced it as a late rejection |

#### Not changed

- `Bun.serve` error reporting for a body that fails after the status
line is committed. `error()` stays bypassed (#4251), the failure prints
to stderr, and the production async path stays quiet.
- `controller.write()` / `flush()` on a native sink still throw the
"already been closed… once `async pull()` returns" message after the
source closed the stream itself. That message is what catches a `pull()`
that forgot to `await`.
- The JS reader path still calls `pull()` again per read (never while a
previous call is pending). `write()` after close returns `0` there, as
before.
- Standalone `Bun.ArrayBufferSink` and `Bun.file().writer()` still throw
after `end()`.
- Async generator and `Response` bodies. New tests pin the existing
behavior.

#### Refactors in this PR (no behavior change intended)

- `JSReadableSinkControllerBase` holds the source cell and the close
promise itself. `detach()`, `JSSinkController__onClose`, and
`JSSinkController__onReady` call
`Bun::WebStreams::sinkControllerOnClose` / `OnReady` directly, and
`end()` dispatches on `SinkID`. The `JSDirectSinkCloseState` cell and
three `JSBoundFunction` handlers are deleted. Every consumer of those
handlers was C++.
- Native stream sources (`Blob` / `File` / `Bytes`) notify their C++
adapter through `Bun__NativeStreamSourceAdapter__onClose` instead of a
bound function stored through a JS setter. `handle.onDrain` had no
caller and is deleted, along with `close_handler` / `close_ctx`.
- `Blob.rs` `FileStreamWrapper` is deleted. `FileSink::pipe_stream`
covers that path.

### How did you verify your code works?

- `test/js/web/streams/direct-stream-contract.ts` is a source-shapes ×
consumers matrix. Each cell asserts that `pull()` ran once, the bytes
delivered, the error, and that `cancel()` ran only on a consumer abort.
- Shapes whose `pull()` returns or resolves without closing are marked
`oneShotOnly` and skip the reader consumers, which would call `pull()`
again.
- Consumers in `streams.test.js`: `.text()` / `.bytes()` / `.blob()`,
`readableStreamTo*`, `getReader()`, `for await`, `pipeTo`,
`pipeThrough`, `tee`, `Bun.write`, `Bun.spawn` stdin.
- Consumers in `serve-direct-readable-stream.test.ts`: `Bun.serve` HTTP
and HTTPS bodies, `fetch` request bodies.
- `serve-direct-readable-stream.test.ts` also covers the `close()` /
`end()` matrix over HTTP, HTTPS, and HTTP/3, the last-chunk orderings,
`close()` under transport backpressure, keep-alive after `close()`, and
client-abort races.
- On 1.4.3, 30 cases in `serve-direct-readable-stream.test.ts` fail: 6
sync `close()` cases over HTTP and HTTPS, 2 over HTTP/3, 7 last-chunk
orderings, 4 backpressure cases, 8 `cancel()` cases, the keep-alive
`close()` test, and 2 `stop(true) after end()` expectations.
- `bun-write.test.js` and `spawn-stdin-readable-stream.test.ts` cover
the `FileSink` truncation and the GC-after-free case.
- Run locally on the debug ASAN build: `streams.test.js`,
`bun-write.test.js`, `spawn-stdin-readable-stream.test.ts`. The
`Bun.serve` suites run in CI.

<details><summary>Notes</summary>

**Folded PRs**

- #37696: `close()` with a buffered tail sends it at once. Same change
as bug 1. Its HTTP / HTTPS matrix and two HTTP/3 cases are in the test
file. Closed in favor of this PR.
- #41893: set `ended_response` where the uWS end happens. Same as bug 2.
Its seven orderings are in the test file, plus a `close(error)` ordering
that pins only the framing. Closed in favor of this PR.

**Two existing expectations change**

- `server.stop(true) from inside pull() after end()` no longer records a
`cancel()` event.
- The `closeThenPark` fixture yields once before `close()`, because a
synchronous `close()` now completes the request inside the dispatch, as
a synchronous `end()` already did.

**Follow-on from that**

When a synchronous `close()` inside an `async pull()` completes the
response before `pull()` first yields, `do_render_stream` takes its
`has_responded()` path and releases the request while `pull()` keeps
running. `end()` already behaved this way. A `pull()` that throws after
that point is an unhandled rejection. Holding the request until `pull()`
settles would keep the context alive for as long as `pull()` runs, so it
is left out.

**Implementation notes**

- `fail()` never sends the buffered tail. The previous `end()` body
parked it with `requested_end` set, and the auto-flusher (which
`do_render_stream`'s `drain_microtasks()` runs) then sent it with
`res.end()`, completing a failed body. A sync `pull(c){ c.write(x);
c.close(err) }` hit that.
- A sync `pull()` that calls `close(error)` makes `readDirectStream`
return a promise rejected with the error (marked handled: every caller
is a native owner that reads its state), so owners take the same reject
path as the async shape.
- `JsSinkType::end` is reached only from `js_close` with an empty
reason, so its `err` argument was always `None` here. `global_this` is
set when `do_render_stream` constructs the sink, before
`assign_to_stream` runs.
- A parked `pending_flush` from a sync `close()` under backpressure is
picked up by `do_render_stream` and `handle_resolve_stream`, the same
way a sync `end()` already worked over HTTP/3.
- #41581, #41721, #41561, and #41757 touch the same lines and will need
a rebase on whichever lands first.

</details>

<!-- robobun:evidence:begin -->

---

**no test proof** · iteration 3 · platform-specific test(s) that do not
run on this machine, deferring to CI, which covers all platforms:
test/js/web/streams/streams.test.js,
test/js/bun/spawn/spawn-stdin-readable-stream.test.ts,
test/js/bun/io/bun-write.test.js

<!-- robobun:evidence:end -->

---------

Co-authored-by: Jarred Sumner <jarred@jarredsumner.com>
Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

ReadbleStream reads to completion in server, does not stream

2 participants