From 8d877b0926964e2018703d9deb6dd3141b2664f5 Mon Sep 17 00:00:00 2001 From: S1ro1 Date: Wed, 25 Mar 2026 18:04:46 +0530 Subject: [PATCH 1/2] Feat: add weigh transfer fix --- src/prime_rl/inference/patches.py | 38 +++++++++++++++++++++++++++++++ 1 file changed, 38 insertions(+) diff --git a/src/prime_rl/inference/patches.py b/src/prime_rl/inference/patches.py index 3295dc5c37..0aac0bb60b 100644 --- a/src/prime_rl/inference/patches.py +++ b/src/prime_rl/inference/patches.py @@ -10,6 +10,7 @@ def transformers_v5_compat(): Qwen3VLMoeTextConfig.tie_word_embeddings = False _patch_qwen35_lora() + monkey_patch_dp_engine_core_pause_resume_deadlock() def _patch_qwen35_lora(): @@ -712,3 +713,40 @@ def _apply_with_saved_kernel(self, layer, x, topk_weights, topk_ids, shared_expe self.base_layer._replace_quant_method(new_method) FusedMoEWithLoRA._inject_lora_into_fused_moe = _fixed_inject + + +def monkey_patch_dp_engine_core_pause_resume_deadlock(): + """Fix deadlock with pause/resume and collective_rpc in DP engine core. + + When a request arrives for an already-completed wave while the scheduler is + paused, the unpatched code sends a start_wave notification that triggers + collective_rpc on other DP engines. But the paused engine can't participate + in the collective, causing a deadlock. + + Fix: only send the start_wave notification when the scheduler is unpaused, + and explicitly set engines_running=True before notifying. + + Upstream: https://github.com/vllm-project/vllm/pull/37024 + """ + from vllm.v1.core.sched.interface import PauseState + from vllm.v1.engine import EngineCoreOutputs + from vllm.v1.engine.core import DPEngineCoreProc, EngineCore + from vllm.v1.request import Request + + _base_add_request = EngineCore.add_request + + def _patched_add_request(self, request: Request, request_wave: int = 0): + _base_add_request(self, request, request_wave) + if self.has_coordinator and request_wave != self.current_wave: + if request_wave > self.current_wave: + self.current_wave = request_wave + elif ( + not self.engines_running + and self.scheduler.pause_state == PauseState.UNPAUSED + ): + self.engines_running = True + self.output_queue.put_nowait( + (-1, EngineCoreOutputs(start_wave=self.current_wave)) + ) + + DPEngineCoreProc.add_request = _patched_add_request From 1c192cc04384e0d1e2271993b564bdcb2b8ab8e5 Mon Sep 17 00:00:00 2001 From: S1ro1 Date: Wed, 25 Mar 2026 22:27:41 +0530 Subject: [PATCH 2/2] Fix style --- src/prime_rl/inference/patches.py | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/src/prime_rl/inference/patches.py b/src/prime_rl/inference/patches.py index 0aac0bb60b..480198e638 100644 --- a/src/prime_rl/inference/patches.py +++ b/src/prime_rl/inference/patches.py @@ -740,13 +740,8 @@ def _patched_add_request(self, request: Request, request_wave: int = 0): if self.has_coordinator and request_wave != self.current_wave: if request_wave > self.current_wave: self.current_wave = request_wave - elif ( - not self.engines_running - and self.scheduler.pause_state == PauseState.UNPAUSED - ): + elif not self.engines_running and self.scheduler.pause_state == PauseState.UNPAUSED: self.engines_running = True - self.output_queue.put_nowait( - (-1, EngineCoreOutputs(start_wave=self.current_wave)) - ) + self.output_queue.put_nowait((-1, EngineCoreOutputs(start_wave=self.current_wave))) DPEngineCoreProc.add_request = _patched_add_request