Skip to content

KTOR-9943 Cancel abandoned streaming responses in HttpStatement - #5913

Open
Osip Fatkullin (osipxd) wants to merge 4 commits into
release/3.xfrom
osipxd/client-stream-cleanup
Open

Osip Fatkullin (osipxd) wants to merge 4 commits into
release/3.xfrom
osipxd/client-stream-cleanup

Conversation

@osipxd

@osipxd Osip Fatkullin (osipxd) commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

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:

  • Resolve delegated responses to their origin before inspecting the engine channel. This avoids calling a delegated rawContent getter that may create a new channel or throw.
  • Propagate response-job cancellation to DefaultTransform's output channel, so its copying writer can stop when the transformed body is abandoned.

It's better to review by commits.

@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Sep 30, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-09-30T16:20:29.608805Z a62da17 New commits
🔒 Security Review ✅ Completed 2026-09-30T16:24:58.751621Z a62da17 New commits
ℹ️ About Codex in GitHub

Your 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.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 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()

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment on lines +247 to +250
return when (val origin = resolveOrigin()) {
is DefaultHttpResponse -> origin.rawContent
is SavedHttpResponse -> ByteReadChannel.Empty
else -> null

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge 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 👍 / 👎.

// 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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Btw, what do you mean by "Reading closure"?

@osipxd Osip Fatkullin (osipxd) Oct 1, 2026 •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ohh, don't mind, I read "closure" in the meaning of "lambda"

responseJobHolder.complete()
}
// response job cancels → cancel output
byteChannel.attachJob(response.coroutineContext.job)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

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.

2 participants