Repository navigation
[core] Clean up failed placement group commits before retry - #66876
hahahahahayesyeseys wants to merge 10 commits into
Conversation
Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
There was a problem hiding this comment.
Code Review
This pull request improves the cleanup of uncommitted placement group bundles before retrying by introducing DestroyPlacementGroupUncommittedBundleResources and updating RemovePlacementGroupBundles to support asynchronous callbacks. It also refactors NodeManager::HandleRemovePlacementGroupBundles to target only the specific bundles being removed rather than the entire placement group. Feedback focuses on improving safety and API design, specifically recommending passing BundleLocations by const reference instead of shared pointer, adding null checks for uncommitted_bundle_locations and the callback to prevent potential crashes, and adding a null check for workers in FakeWorkerPool.
| void GcsPlacementGroupScheduler::DestroyPlacementGroupUncommittedBundleResources( | ||
| const PlacementGroupID &placement_group_id, | ||
| const std::shared_ptr<BundleLocations> &bundle_locations, | ||
| rpc::StatusCallback callback) { | ||
| auto bundles_per_node = GroupBundlesByNode(*bundle_locations); | ||
| if (bundles_per_node.empty()) { | ||
| callback(Status::OK()); | ||
| return; | ||
| } | ||
|
|
||
| auto pending_nodes = std::make_shared<std::atomic<size_t>>(bundles_per_node.size()); | ||
| auto cleanup_failed = std::make_shared<std::atomic<bool>>(false); | ||
| for (auto &entry : bundles_per_node) { | ||
| RemovePlacementGroupBundles( | ||
| placement_group_id, | ||
| entry.second, | ||
| gcs_node_manager_.GetAliveNode(entry.first), | ||
| /*max_retry*/ 5, | ||
| /*current_retry_count*/ 0, | ||
| [pending_nodes, cleanup_failed, callback](const Status &status) { | ||
| if (!status.ok()) { | ||
| cleanup_failed->store(true); | ||
| } | ||
| if (pending_nodes->fetch_sub(1) == 1) { | ||
| callback(cleanup_failed->load() | ||
| ? Status::IOError("Failed to remove one or more bundle groups.") | ||
| : Status::OK()); | ||
| } | ||
| }); | ||
| } | ||
| } |
There was a problem hiding this comment.
Update the implementation of DestroyPlacementGroupUncommittedBundleResources to accept const BundleLocations & directly. Additionally, check if callback is valid before invoking it to prevent potential std::bad_function_call crashes when called with an empty callback.
void GcsPlacementGroupScheduler::DestroyPlacementGroupUncommittedBundleResources(
const PlacementGroupID &placement_group_id,
const BundleLocations &bundle_locations,
rpc::StatusCallback callback) {
auto bundles_per_node = GroupBundlesByNode(bundle_locations);
if (bundles_per_node.empty()) {
if (callback) {
callback(Status::OK());
}
return;
}
auto pending_nodes = std::make_shared<std::atomic<size_t>>(bundles_per_node.size());
auto cleanup_failed = std::make_shared<std::atomic<bool>>(false);
for (auto &entry : bundles_per_node) {
RemovePlacementGroupBundles(
placement_group_id,
entry.second,
gcs_node_manager_.GetAliveNode(entry.first),
/*max_retry*/ 5,
/*current_retry_count*/ 0,
[pending_nodes, cleanup_failed, callback](const Status &status) {
if (!status.ok()) {
cleanup_failed->store(true);
}
if (pending_nodes->fetch_sub(1) == 1) {
if (callback) {
callback(cleanup_failed->load()
? Status::IOError("Failed to remove one or more bundle groups.")
: Status::OK());
}
}
});
}
}There was a problem hiding this comment.
Partially applied in 1142c6a: bundle_locations is now passed as const BundleLocations &. I kept the callback required and invoked unconditionally because this is a private cleanup helper and every call site supplies the completion callback that releases the lease tracker and advances scheduling state. Treating an empty callback as normal would leave that state transition incomplete and hide a programming error.
| void DestroyPlacementGroupUncommittedBundleResources( | ||
| const PlacementGroupID &placement_group_id, | ||
| const std::shared_ptr<BundleLocations> &bundle_locations, | ||
| rpc::StatusCallback callback); |
There was a problem hiding this comment.
The DestroyPlacementGroupUncommittedBundleResources function only reads from bundle_locations and does not need to share or take ownership of the shared pointer. Passing it as const BundleLocations & is cleaner, more idiomatic, and avoids any potential null pointer dereferencing issues inside the function.
| void DestroyPlacementGroupUncommittedBundleResources( | |
| const PlacementGroupID &placement_group_id, | |
| const std::shared_ptr<BundleLocations> &bundle_locations, | |
| rpc::StatusCallback callback); | |
| void DestroyPlacementGroupUncommittedBundleResources( | |
| const PlacementGroupID &placement_group_id, | |
| const BundleLocations &bundle_locations, | |
| rpc::StatusCallback callback); |
There was a problem hiding this comment.
Applied in 1142c6a: the helper now accepts const BundleLocations &, and both call sites dereference the tracker-owned container explicitly.
| const auto &uncommitted_bundle_locations = | ||
| lease_status_tracker->GetUnCommittedBundleLocations(); | ||
| auto committed_bundle_locations = std::make_shared<BundleLocations>(); | ||
| for (const auto &[bundle_id, location] : *prepared_bundle_locations) { | ||
| if (!uncommitted_bundle_locations->contains(bundle_id)) { | ||
| committed_bundle_locations->emplace(bundle_id, location); | ||
| } | ||
| } |
There was a problem hiding this comment.
To prevent potential null pointer dereferences, check if uncommitted_bundle_locations is valid before dereferencing it. Additionally, if it is null or empty, we can directly assign prepared_bundle_locations to committed_bundle_locations to avoid unnecessary copying.
const auto &uncommitted_bundle_locations =
lease_status_tracker->GetUnCommittedBundleLocations();
auto committed_bundle_locations = std::make_shared<BundleLocations>();
if (uncommitted_bundle_locations) {
for (const auto &[bundle_id, location] : *prepared_bundle_locations) {
if (!uncommitted_bundle_locations->contains(bundle_id)) {
committed_bundle_locations->emplace(bundle_id, location);
}
}
} else {
committed_bundle_locations = prepared_bundle_locations;
}There was a problem hiding this comment.
I kept the existing invariant rather than adding a nullable fallback. LeaseStatusTracker constructs and owns uncommitted_bundle_locations_, and the accessor always returns that initialized container. Treating null as all-committed would both hide an invalid tracker state and risk recording failed bundles as committed.
| DestroyPlacementGroupUncommittedBundleResources( | ||
| placement_group_id, | ||
| uncommitted_bundle_locations, | ||
| [placement_group, schedule_failure_handler](const Status &status) { | ||
| if (!status.ok()) { | ||
| RAY_LOG(ERROR) << "Failed to clean up uncommitted placement group bundles: " | ||
| << status; | ||
| } | ||
| schedule_failure_handler(placement_group, /*is_feasible*/ true); | ||
| }); |
There was a problem hiding this comment.
Update the call site to check if uncommitted_bundle_locations is non-null and dereference it to match the updated signature of DestroyPlacementGroupUncommittedBundleResources.
| DestroyPlacementGroupUncommittedBundleResources( | |
| placement_group_id, | |
| uncommitted_bundle_locations, | |
| [placement_group, schedule_failure_handler](const Status &status) { | |
| if (!status.ok()) { | |
| RAY_LOG(ERROR) << "Failed to clean up uncommitted placement group bundles: " | |
| << status; | |
| } | |
| schedule_failure_handler(placement_group, /*is_feasible*/ true); | |
| }); | |
| if (uncommitted_bundle_locations) { | |
| DestroyPlacementGroupUncommittedBundleResources( | |
| placement_group_id, | |
| *uncommitted_bundle_locations, | |
| [placement_group, schedule_failure_handler](const Status &status) { | |
| if (!status.ok()) { | |
| RAY_LOG(ERROR) << "Failed to clean up uncommitted placement group bundles: " | |
| << status; | |
| } | |
| schedule_failure_handler(placement_group, /*is_feasible*/ true); | |
| }); | |
| } else { | |
| schedule_failure_handler(placement_group, /*is_feasible*/ true); | |
| } |
There was a problem hiding this comment.
The call site now passes *uncommitted_bundle_locations for the const-reference signature. I did not add a null branch because the container is initialized by every LeaseStatusTracker constructor; null is not a valid state and should not silently bypass cleanup.
| DestroyPlacementGroupUncommittedBundleResources( | ||
| placement_group_id, | ||
| uncommitted_bundle_locations, | ||
| [placement_group, schedule_failure_handler](const Status &status) { | ||
| if (!status.ok()) { | ||
| RAY_LOG(ERROR) << "Failed to clean up uncommitted placement group bundles: " | ||
| << status; | ||
| } | ||
| schedule_failure_handler(placement_group, /*is_feasible*/ true); | ||
| }); |
There was a problem hiding this comment.
Update the call site to check if uncommitted_bundle_locations is non-null and dereference it to match the updated signature of DestroyPlacementGroupUncommittedBundleResources.
if (uncommitted_bundle_locations) {
DestroyPlacementGroupUncommittedBundleResources(
placement_group_id,
*uncommitted_bundle_locations,
[placement_group, schedule_failure_handler](const Status &status) {
if (!status.ok()) {
RAY_LOG(ERROR) << "Failed to clean up uncommitted placement group bundles: "
<< status;
}
schedule_failure_handler(placement_group, /*is_feasible*/ true);
});
} else {
schedule_failure_handler(placement_group, /*is_feasible*/ true);
}There was a problem hiding this comment.
The call site now passes *uncommitted_bundle_locations for the const-reference signature. I did not add a null branch because the container is initialized by every LeaseStatusTracker constructor; null is not a valid state and should not silently bypass cleanup.
| for (const auto &worker : registered_workers) { | ||
| if (worker->Connection() == connection) { | ||
| return worker; | ||
| } | ||
| } |
There was a problem hiding this comment.
Add a null check for worker in the loop to prevent potential null pointer dereferences if registered_workers contains any null shared pointers.
| for (const auto &worker : registered_workers) { | |
| if (worker->Connection() == connection) { | |
| return worker; | |
| } | |
| } | |
| for (const auto &worker : registered_workers) { | |
| if (worker && worker->Connection() == connection) { | |
| return worker; | |
| } | |
| } |
There was a problem hiding this comment.
I left this unchanged. registered_workers is a test-owned collection populated with registered FakeWorker instances; a null entry is not a valid worker-pool state. Silently skipping one would mask a malformed fixture rather than protect a reachable production path.
Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
…mit-bundle-index Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
…mit-bundle-index Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
…mit-bundle-index Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
…mit-bundle-index Signed-off-by: hahahahahayesyeseys <1145018948@qq.com>
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
Reviewed by Cursor Bugbot for commit e019f1b. Configure here.
| // A queued wildcard lease is not tied to a concrete bundle yet. Cancel it | ||
| // conservatively so it cannot remain queued after placement group cleanup. | ||
| return true; | ||
| } |
There was a problem hiding this comment.
Cleanup cancels live wildcard leases
Medium Severity
Failed-commit cleanup now sends RemovePlacementGroupBundles while the placement group is still alive and will be retried, but unallocated wildcard leases are cancelled with SCHEDULING_CANCELLED_PLACEMENT_GROUP_REMOVED. Clients treat that as a permanent placement-group removal, so queued wildcard tasks and actors fail even though remaining bundles can still serve them or the commit will be retried.
Additional Locations (1)
Reviewed by Cursor Bugbot for commit e019f1b. Configure here.


Description
When
CommitBundleResourcesfailed on a live raylet, the failed bundle was still added tocommitted_bundle_location_index_. A later retry could place it on another node while the reverse index continued to point to the old node, so removing the placement group leaked the reservation on the new node.This change:
RemovePlacementGroupBundlescancel leases, stop workers, and return resources only for the bundles named in the request, preserving other bundles from the same placement group;bundle_index=-1workers through their actual non-null indexed placement-group resources, so deleting one bundle stops wildcard workers allocated to that bundle without stopping wildcard workers allocated to surviving bundles;The regression tests cover the full failure lifecycle: commit RPC failure, concurrent cancellation during cleanup, cleanup on the old node, retry on a new node, correct reverse-index ownership, final placement-group removal, worker isolation between two bundles in one placement group, allocated and unallocated wildcard cleanup, and retry after held resources are released.
Related issues
Fixes #66873
Additional information
Validation
The final pushed head
e019f1b56be77d8fef406fe3571c471d1671bdf6, synchronized withmasterat990ff9f6ff1977516e866aca29d7313000b83efe, passed all affected C++ targets:On that same head, each of these bundle-removal boundaries passed 20 runs:
RemovePlacementGroupBundlesOnlyStopsRequestedBundleWorkers,RemovePlacementGroupBundlesCancelsUnallocatedWildcardLease, andRemovePlacementGroupBundlesReturnsRetryableFailure. They jointly verify that removing bundle 0 stops its wildcard worker, preserves a wildcard worker allocated to surviving bundle 1 even when its other allocation slot is null, still cancels an unallocated wildcard lease, and reports in-use resources as retryable. Earlier synchronized heads also passed 20-run stress forFailedCommitIsCleanedBeforeRetryandFailedCommitCanBeCancelledDuringCleanup.Ray's applicable pre-commit checks passed for the changed C++ files, including
cpplint,clang-format, link validation, and the C++ include check. The complete seven-file PR diff passedgit diff --check, and every PR commit has a DCO sign-off.Before the implementation, the first GCS regression failed because no cleanup request was sent after the commit error, and the Raylet regression failed because removing one bundle also removed the other bundle's worker. Before the review follow-up, the cancellation regression triggered the reported
RAY_CHECKbecause the lease tracker had already been removed while cleanup was still pending. Buildkite microcheck #55906 then exposed a wildcard worker whose logical bundle index remained-1while its allocated resources belonged to bundle 0; cleanup missed that worker and crashed the raylet whenReturnBundlereported that resources were still in use. The updated regression also reproduces the later review boundary where one of a granted wildcard worker's two allocation slots is null: null is no longer treated as a removed-bundle match, while the surviving bundle's real allocation is honored.Buildkite microcheck #55920 also showed that two existing Python observability tests rely on the stable
placement group was removedsubstring in worker exit details. The updated detail preserves that substring while distinguishing removal of an assigned bundle. Those Python consumers could not be functionally executed locally: the Bazel test environment initially lackedpytest, and after supplying a compatiblepytest, it lacked an installed Ray package.The failing
placement_group_exampleintegration target could not enter Ray locally because this checkout does not contain a builtray._raylet, and Docker access is unavailable to the current user. The push-triggered Buildkite run remains the final integration check for that target and the two Python consumers.Duplicate-work check
As of 2026-10-09, issue #66873 is open, unassigned, and has no linked branch or pull request. No open PR references #66873 or addresses live-node commit failures, stale committed-bundle indexing, and per-bundle Raylet cleanup together. #66188 touches the same scheduler but fixes cancellation during the prepare phase and scheduling-token release for #64693; it does not change commit-error indexing or Raylet bundle-removal scope. This PR is the scoped follow-up described in #66769.
Risk and compatibility
There is no public API, configuration, or wire-format change. The behavioral change is limited to placement-group commit-failure cleanup and the existing bundle-removal RPC's server-side scope.
Screenshots are not applicable because this changes internal C++ scheduling and resource-cleanup behavior.
AI assistance disclosure
AI assistance was used to prepare the implementation, tests, and this description.