[fix](ray) Publish completion after stream retirement - #237
Merged
kaka11chen merged 1 commit intoJul 28, 2026
Merged
Conversation
Contributor
Author
|
@codex review |
|
Codex Review: Didn't find any major issues. 🚀 Reviewed commit: ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
If Codex has suggestions, it will comment; otherwise it will react with 👍. Codex can also answer questions or update the PR. Try commenting "@codex address that feedback". |
hubgeter
approved these changes
Jul 28, 2026
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.
Summary
RayStreamAdapter.retire()releases task/source ownership and deletes the ObjectRef streamcompleteevent only after that local retirement finishesRoot cause
AsyncResultCollector._maybe_complete_record()removed the record and publishedcompletebefore callingrecord.adapter.retire(). Consumers could therefore observe an empty collector while the adapter still held strong references toTaskLeaseObjectRefGenerator, making the weakref assertion scheduler-dependent in the full fast suite.Testing
python -m pytest tests/fast/test_udf_stream_result_collector.py(21 passed)python -m pytest tests/fast/test_udf_process.py(6 passed)scripts/run_release_tests.sh(418 passed, 1 skipped; 7 real-Ray tests passed)scripts/run_fast_tests.sh non-ray(5895 passed, 28 skipped, 8 xfailed, 1 xpassed)pre-commit run --from-ref upstream/main --to-ref HEADCloses #217