Repository navigation
[TSL] Wait for queued and running work in default UnboundedWorkQueue::~UnboundedWorkQueue - #3653
Merged
Merged
Conversation
copybara-service
Bot
force-pushed
the
test_993490972
branch
from
October 5, 2026 08:31
be8cdd4 to
208bb75
Compare
…::~UnboundedWorkQueue`
In the default/OSS implementation of `tsl::UnboundedWorkQueue` (`xla/tsl/platform/default/unbounded_work_queue.{h,cc}`), `~UnboundedWorkQueue()` previously set `cancelled_ = true` immediately without waiting for pending or in-flight work items to finish, and then acquired `thread_pool_mu_` while joining `thread_pool_` threads.
This differed from the Google-internal implementation (which waits for all work to complete in `~UnboundedWorkQueue()`) and caused a lock-order inversion deadlock when an in-flight worker thread called `UnboundedWorkQueue::Schedule()` concurrently with `~UnboundedWorkQueue()` (e.g., in `//xla/pjrt/c:pjrt_c_api_gpu_test_nvgpu_any` during `PjrtCApiTest.BufferTransferImmutableUntilTransferCompletes`, where `CopyRawHostToDeviceAndReturnEvent` calls `ExecuteWhenReady` -> `Schedule()` on the `H2D Dispatch` worker thread after `on_done_with_host_buffer` has already unblocked client destruction):
- Main thread in `~UnboundedWorkQueue()` holds `thread_pool_mu_` and blocks in `pthread_join` waiting for the worker thread.
- Worker thread in `Schedule()` holds `work_queue_mu_` and blocks trying to acquire `thread_pool_mu_`.
Update `default/unbounded_work_queue.{h,cc}` to track `num_running_functions_` (guarded by `work_queue_mu_`) and wait in `~UnboundedWorkQueue()` until `work_queue_.empty() && num_running_functions_ == 0` before setting `cancelled_ = true` and joining `thread_pool_`. Also update `unbounded_work_queue_test.cc` to expect all closures to finish in `RacyDestructor` and add `NestedClosureDuringDestructor` to cover scheduling nested work during destruction.
PiperOrigin-RevId: 993570944
copybara-service
Bot
force-pushed
the
test_993490972
branch
from
October 5, 2026 10:10
208bb75 to
32b347f
Compare
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.
[TSL] Wait for queued and running work in default
UnboundedWorkQueue::~UnboundedWorkQueueIn the default/OSS implementation of
tsl::UnboundedWorkQueue(xla/tsl/platform/default/unbounded_work_queue.{h,cc}),~UnboundedWorkQueue()previously setcancelled_ = trueimmediately without waiting for pending or in-flight work items to finish, and then acquiredthread_pool_mu_while joiningthread_pool_threads.This differed from the Google-internal implementation (which waits for all work to complete in
~UnboundedWorkQueue()) and caused a lock-order inversion deadlock when an in-flight worker thread calledUnboundedWorkQueue::Schedule()concurrently with~UnboundedWorkQueue()(e.g., in//xla/pjrt/c:pjrt_c_api_gpu_test_nvgpu_anyduringPjrtCApiTest.BufferTransferImmutableUntilTransferCompletes, whereCopyRawHostToDeviceAndReturnEventcallsExecuteWhenReady->Schedule()on theH2D Dispatchworker thread afteron_done_with_host_bufferhas already unblocked client destruction):~UnboundedWorkQueue()holdsthread_pool_mu_and blocks inpthread_joinwaiting for the worker thread.Schedule()holdswork_queue_mu_and blocks trying to acquirethread_pool_mu_.Update
default/unbounded_work_queue.{h,cc}to tracknum_running_functions_(guarded bywork_queue_mu_) and wait in~UnboundedWorkQueue()untilwork_queue_.empty() && num_running_functions_ == 0before settingcancelled_ = trueand joiningthread_pool_. Also updateunbounded_work_queue_test.ccto expect all closures to finish inRacyDestructorand addNestedClosureDuringDestructorto cover scheduling nested work during destruction.