KTOR-9943 Cancel abandoned streaming responses in HttpStatement - #5913
Osip Fatkullin (osipxd) wants to merge 4 commits into
Conversation
0ffcf71 to
f18dd59
Compare
f18dd59 to
5a9abf7
Compare
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
5a9abf7 to
c4c8582
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 5a9abf7648
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "Codex (@codex) review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "Codex (@codex) address that feedback".
| join() | ||
| when (cause) { | ||
| is CancellationException -> job.cancel(cause) | ||
| null if content?.isClosedForRead == true && content.closedCause == null -> job.complete() |
There was a problem hiding this comment.
Cancel when transformed output remains unread after engine EOF
When a caller partially consumes bodyAsChannel() and exits after the copier has read the final engine bytes but while it is blocked writing unread output, the engine channel is already closed, so this branch completes the response job. The responseJobHolder created in DefaultTransform.kt remains a child until that blocked writer finishes, while attachJob only cancels its output if the response job is cancelled; consequently job.join() can wait forever. The same pattern applies to decompression writers. Treating engine EOF alone as proof that all transformed output was consumed therefore reintroduces the abandoned-response hang.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Yes, this problem is tracked here: KTOR-9953 HttpStatement.execute hangs when the engine body reaches EOF but bodyAsChannel output is full
I want to fix it in a separate PR.
| return when (val origin = resolveOrigin()) { | ||
| is DefaultHttpResponse -> origin.rawContent | ||
| is SavedHttpResponse -> ByteReadChannel.Empty | ||
| else -> null |
There was a problem hiding this comment.
Preserve successful completion for custom response wrappers
When a public receive-pipeline interceptor returns its own HttpResponse wrapper, as the repository's ClientPipelinesTest already demonstrates, resolveOrigin() leaves that wrapper unchanged and this fallback returns null. Cleanup then cancels the response job even if the wrapper's raw channel was fully consumed successfully, causing completion handlers and returned job references to observe a cancellation rather than normal completion. Unknown wrappers with accessible content need a completion check instead of being classified unconditionally as abandoned.
Useful? React with 👍 / 👎.
c4c8582 to
a62da17
Compare
| // Pre-read closedCause to access the getter only once | ||
| // Reading closure first avoids mistaking a cancellation between the reads for clean EOF | ||
| val contentIsClosedForRead = content?.isClosedForRead == true | ||
| val contentCloseCause = content?.closedCause |
There was a problem hiding this comment.
I guess we should avoid getters for functions that can allocate new objects in the future
FYI, I've checked the implementation of ByteChannel.isClosedForRead, and it accesses closedCause, which can lead to one more allocation
| } | ||
| if (cause == null) { | ||
| // Pre-read closedCause to access the getter only once | ||
| // Reading closure first avoids mistaking a cancellation between the reads for clean EOF |
There was a problem hiding this comment.
Btw, what do you mean by "Reading closure"?
There was a problem hiding this comment.
It's about content?.isClosedForRead. We should read it before reading closedClosure, because otherwise we can get isClosedForRead == false and closedClosure != null
I should rewrite this comment to make it clear
There was a problem hiding this comment.
Ohh, don't mind, I read "closure" in the meaning of "lambda"
| responseJobHolder.complete() | ||
| } | ||
| // response job cancels → cancel output | ||
| byteChannel.attachJob(response.coroutineContext.job) |
There was a problem hiding this comment.
Just noticed that attachJob uses invokeOnCompletion(onCancelling = true) and KDoc says
Kept for preserving compatibility. Shouldn't be used by anyone.
Is there any reason for it?
There was a problem hiding this comment.
IIRC onCancelling = true makes it fire at the moment state is changed to Cancelling without waiting for cancellation finishing (Cancelled).
It's an option 1 from here.
But I agree that ideally we shouldn't use internal coroutine APIs.


Subsystem
Client
Motivation
KTOR-9943 HttpStatement.execute can hang when its block abandons a response body
HttpStatement.execute {}could hang after the block returned with an unread streaming body. Cleanup completed the response job and waited for it, even though the engine was still handling the abandoned body.Solution
Cleanup now completes the job only when the engine body reached EOF, otherwise it cancels the job. Fully consumed and saved responses retain normal completion.
Also:
rawContentgetter that may create a new channel or throw.DefaultTransform's output channel, so its copying writer can stop when the transformed body is abandoned.It's better to review by commits.