server: decode streaming request bodies in the handler's task - #299
Merged
Merged
Conversation
iainmcgin
marked this pull request as ready for review
September 11, 2026 22:24
A client-streaming or bidi handler received its request messages from a reader task spawned for each call, which decoded envelopes from the HTTP body and passed each message through an `mpsc::channel(1)`. Every message therefore crossed one more task boundary on its way from hyper to the handler. The request stream now owns the body and the `EnvelopeDecoder` and decodes in `poll_next`, in whichever task polls it. The detached reader task, its channel and the `StreamingResponseBody` field that held its join handle are gone. The drain from #313 is kept, with the same bounds: when the stream ends (END_STREAM, a decode error) or is dropped before the body has ended, it drops the decoder with any partial message and hands the rest of the body to a detached task that reads and discards frames until the body ends, 1 MiB has been discarded (counting what the decoder left of the last frame) or `DRAIN_TIMEOUT` has passed. Reading rather than dropping still matters for h2's small-frame budget and for HTTP/1.1 connection reuse. A body that already reports its end needs no task. The drain now starts when the handler drops the stream rather than when the handler's end of a channel is seen to close, so there is nothing to race while a client is stalled. The body is now read only while the handler polls its request stream; the old reader read up to two messages ahead. The drain task is spawned on the runtime the stream was created on, normally the server's, so a handler that drops its request stream on another thread or runtime still gets the drain. `spawn_detached`, the fallback for a stream created outside a runtime, no longer returns a handle and does nothing outside a Tokio runtime instead of panicking. The reader tests now drive `decode_request_body` directly and wait for the body to be dropped instead of joining the reader task; a per-thread panic count stands in for the join handle's panic check. New tests cover frames that split or pack envelopes, trailers and empty frames, a body already at its end, a drop on another thread or outside a runtime, trailing data in the END_STREAM frame that alone exceeds the drain limit, and a drain that stops at a body error. `http1_connection_survives_unread_streaming_request_body` (tests/streaming) checks over a raw socket that a connection whose streaming request was abandoned after a decode error serves the next request; it fails with the drain disabled, as do the four `*_early_return_*` server tests. Stale mentions of a body reader in the `intercept_head` docs and the guide are updated. Fixes #298. Signed-off-by: Iain McGinniss <309153+iainmcgin@users.noreply.github.com>
iainmcgin
force-pushed
the
iain/inline-request-body-reader
branch
from
September 23, 2026 02:30
c398850 to
d1f7bd3
Compare
iainmcgin
marked this pull request as draft
September 23, 2026 02:30
Signed-off-by: Iain McGinniss <309153+iainmcgin@users.noreply.github.com>
iainmcgin
marked this pull request as ready for review
September 23, 2026 19:49
rpb-ant
approved these changes
Sep 23, 2026
iainmcgin
added a commit
that referenced
this pull request
Sep 26, 2026
…299) Re-runs the latency, echo, log-ingest and client-stack benchmarks on a bare-metal c7i.metal-24xl (turbo disabled, client and servers on one host over loopback, unpinned) and rewrites those tables, the charts and the number-tied prose from that run. The fortunes and CPU-profile subsections keep their 2026-03 numbers, and the section header now says that the 2026-09 date applies unless a subsection says otherwise. The 10-message client stream is now level across the three stacks (166.3 us for connectrpc-rs against 168.2 us for tonic and 162.2 us for tonic-protobuf), where the previous run had connectrpc-rs 9-14% slower; the server now decodes streamed request messages in the handler's task, so the summary drops the cross-task explanation. The log-batch figures are now stated from the table's baseline (tonic 6-13% fewer requests per second under load, tonic-protobuf within 4%), and the latency figure as time taken (39% longer with prost, 16% with upb). A second echo and log-ingest pass in the same session reproduced every echo cell within 1.1% and every log-ingest cell within 1.4%, except tonic-protobuf at c=256, which measured 23% lower on the second pass; the README says so. The client-stacks prose no longer calls grpc-rust's channel non-hyper, since its transport is built on hyper too, and the cross bench's Go toolchain requirement is stated. Signed-off-by: Iain McGinniss <309153+iainmcgin@users.noreply.github.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
The request stream handed to a client-streaming or bidi handler now owns the HTTP body and the envelope decoder and decodes in
poll_next, in the handler's task. The per-call reader task and itsmpscchannel are gone, and with them the cross-task hand-off per message described in #298. Measured on twoc7i.metal-24xlhosts (turbo off, a 10-message gRPC client stream, medians of 5 alternating 10 s rounds), with the server on 4 tokio workers pinned to 4 cores:mainThe
futex, eventfdwriteand extraepoll_waitcalls onmainare the hand-offs between workers. With a single worker the gap shrinks to 9% of server CPU per stream at 1 in flight.When the stream reaches END_STREAM or a decode error, or the handler drops it (returned, interceptor rejection, request timeout), the partial message is freed at once and the unread rest of the body is read and discarded on a detached task, for at most
MAX_DRAIN_BYTES(1 MiB) andDRAIN_TIMEOUT(5 s). The drain is spawned on the runtime captured when the stream was created, so a handler that drops its stream on another thread still drains; a stream dropped where no runtime was captured skips the drain instead of panicking.Behaviour change: the body is read only while the handler polls its stream; the reader task read up to two messages ahead. A bidi handler that keeps its request stream inside its response stream starts draining only when the response body is dropped. Until then hyper's HTTP/1.1 dispatcher reads at most one more chunk of the request body. The response is still written, but a client that must finish sending before it reads waits until the handler drops its stream.
Fixes #298.