Skip to content

feat: add ResponseBodyControl.execute to decide in sequence with body delivery - #2359

Open
mkurz wants to merge 1 commit into
AsyncHttpClient:mainfrom
mkurz:feature/response-body-control-execute
Open

mkurz wants to merge 1 commit into
AsyncHttpClient:mainfrom
mkurz:feature/response-body-control-execute

Conversation

@mkurz

@mkurz mkurz commented Oct 7, 2026

Copy link
Copy Markdown
Contributor

Summary

  • Add ResponseBodyControl.execute(Runnable), which runs a task on the event loop that delivers the response body, in sequence with the callbacks AHC makes there, such as onBodyPartReceived. Calls to suspend(), resume() and cancel() from the task take effect immediately.
  • This lets a consumer that buffers body parts on another thread decide whether to resume based on everything delivered so far, instead of on a decision that a part delivered in the meantime made stale.
  • NettyResponseBodyControl already applied its own control calls this way; that method is now public. The default implementation runs the task on the calling thread.

Problem

A typical streaming consumer buffers body parts from onBodyPartReceived and hands them to a reader on another thread. When that reader has emptied the buffer and still wants more, it calls resume(). Called from that thread, resume() only takes effect later, on the event loop. In between, the event loop can still deliver parts from the current socket read, and those parts may already meet the reader's demand. The queued resume() then lets one more socket read through that nobody asked for.

With a highly compressible body, that one extra read can inflate to many megabytes beyond demand. For a gzipped stream with early demand from another thread, Play WS measured about 35 MB buffered beyond demand, against about 1.9 MB when the decision is re-checked on the event loop.

AHC cannot fix this by re-checking the decision itself, because only the consumer knows which state its decision depends on. Play WS currently works around it by reaching the event loop through Netty's internal ThreadExecutorMap, and then making its resume decision there.

Change

default void execute(Runnable task) {
    task.run();
}
  • On the event loop, for example from within a body callback, the task runs before execute returns.
  • From another thread, it is queued to the event loop and runs there after the work already pending on it. The event loop can deliver further body parts before the task runs; the task sees all of them.
  • Either way, suspend(), resume() and cancel() called from the task take effect immediately.
  • Not every callback runs on the event loop. onThrowable can run on the thread that cancels the request, or on a timer thread when a timeout fires. A task submitted from there is queued like one from any other thread, and can run while that callback is still running. The Javadoc says so; the implementation does not try to change it.
  • The task must not block. Unlike calls to the other methods, it also runs after the response has completed (the interface Javadoc now says so). If the event loop no longer accepts tasks, it does not run.
  • The default implementation runs the task on the calling thread. It suits implementations that apply control calls synchronously on the calling thread, and does not order the task with callbacks on other threads.

NettyResponseBodyControl.execute is the existing private method made public, with a null check. It works for HTTP/1.1 and for HTTP/2 stream channels.

Alternatives considered

  • resumeIf(BooleanSupplier), a resume whose condition AHC evaluates on the event loop. It needs subtle ordering rules: a true condition counts as a resume at that moment and overrides an earlier suspend. A consumer whose own drain loop may run on either thread then has to fit its state machine to those rules, and we found that hard to get right. With execute, such a consumer can write control.execute(() -> { if (cond) control.resume(); }).
  • Demand counting in AHC (request(n)). HTTP/1.1 decodes one socket read into several parts regardless of demand, so AHC would need its own buffer. That's a much larger change.

execute only exposes what the control already does internally. Callbacks already run on the event loop under the same rule not to block.

Compatibility

  • A new default method on ResponseBodyControl. Revapi passed. Existing implementations keep compiling and behave as before, with two exceptions:
    • an implementation that declares a non-public execute(Runnable) no longer compiles;
    • one that already has a public void execute(Runnable) now implements this method, with whatever semantics it has.
  • No behavior change for existing callers.

AI disclosure

Claude Code on behalf of Matthias Kurz. The commit includes Co-Authored-By per AGENTS.md.

Test plan

New tests:

  • ResponseBodyControlTest:
    • execute runs inline in a callback, runs on the event loop when called from another thread, and runs after control calls queued before it;
    • a consumer thread that decides inside execute sees the part delivered in the meantime and does not resume;
    • execute after the response has completed still runs the task on the event loop;
    • the default implementation runs the task on the calling thread.
  • Http2ResponseBodyControlTest: a task from another thread resumes a suspended HTTP/2 stream.
  • NettyResponseBodyControlExecuteTest: once the event loop has shut down, execute neither runs the task nor throws.

Broken variants, each run against these tests and the existing control tests (24 tests in all):

  • The task always runs on the calling thread: 9 failures, including "a task from another thread must run on the channel event loop".
  • The task is always queued: 3 failures, including "a task from a callback must run before execute returns".
  • A rejected task throws: the shutdown test fails.
  • No task after completion: the completion test times out (plus failures of existing tests, since the control's own calls use the same method).

Earlier, at the first revision, deciding on the other thread without execute failed "the task must see the part buffered after it was queued".

Runs:

  • ResponseBodyControlTest, Http2ResponseBodyControlTest and NettyResponseBodyControlExecuteTest on JDK 11: 24 tests passed.

  • ./mvnw -B -ntp -Dgpg.skip=true clean verify on JDK 11: 1,797 tests, 0 failures, 0 errors, 22 skipped; Revapi passed. No test-skipping flags were used.

  • The same 24 tests, compiled on JDK 11 and run on JDK 17, 21 and 25 (Maven Surefire's -Djvm): 0 failures on each.

  • Play WS, with its resume check moved from ThreadExecutorMap to execute, against a local build of the first revision of this branch alone (Scala 2.13):

    • all Play WS tests pass, except a new cookie test that needs the separate raw-Cookie fix;
    • its gzip streaming test buffers 1,875,916 bytes beyond demand, as with the ThreadExecutorMap workaround;
    • with Play WS resuming directly instead of through execute, the same test buffers 35,589,292 bytes and fails.

    This revision changes only Javadoc and tests.

  • All four open AHC changes merged together (raw Cookie headers on retries, ResponseBodyControl.execute, suspend/resume order, demand-bounded decompression): they merge without conflicts, together and in every pair, and ./mvnw -B -ntp -Dgpg.skip=true clean install on JDK 11 passes with 1,853 tests, 0 failures, 0 errors, 22 skipped; Revapi passed.

  • Play WS against that merged build, Scala 2.13 and 3.3, JDK 17: all tests pass.

A consumer that buffers body parts on another thread decides there
whether to resume reads, but resume() from that thread only takes
effect later on the event loop. In between, the event loop can deliver
another part that already meets the consumer's demand, and the resume
then lets one more read through. With a highly compressible body, that
read can inflate to many megabytes beyond demand. Only the caller knows
which state its decision depends on, so AHC cannot re-check it.

Add execute(Runnable), which runs a task on the event loop that
delivers the response body, in sequence with the callbacks made there:
before returning when called on that event loop, and queued to it
otherwise. suspend(), resume() and cancel() called from the task take
effect immediately, so a consumer can decide there on everything
delivered so far. onThrowable can run on other threads, after a cancel
or a timeout, so a task is not ordered with it. Unlike the other
methods, the task also runs after completion; it does not run once the
event loop rejects tasks. NettyResponseBodyControl already applied its
own control calls this way; that method is now public. The default
implementation runs the task on the calling thread, for
implementations that apply control calls synchronously.

Play WS needs this: it currently reaches the event loop through Netty's
internal ThreadExecutorMap to avoid the extra read.

Claude Code on behalf of Matthias Kurz

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.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.

1 participant