From d1acd62d29aa2634d6e19dcdb3d0a2823723f7b3 Mon Sep 17 00:00:00 2001 From: ybyang <10629930+whybeyoung@users.noreply.github.com> Date: Mon, 18 May 2026 22:57:22 +0800 Subject: [PATCH] fix(disagg): unstuck decode aborts under prealloc pressure (#25561) --- python/sglang/srt/disaggregation/decode.py | 6 +++++- python/sglang/srt/managers/scheduler.py | 4 ++++ python/sglang/srt/managers/tokenizer_manager.py | 10 +++++++++- 3 files changed, 18 insertions(+), 2 deletions(-) diff --git a/python/sglang/srt/disaggregation/decode.py b/python/sglang/srt/disaggregation/decode.py index dfa11f492..5be837eb3 100644 --- a/python/sglang/srt/disaggregation/decode.py +++ b/python/sglang/srt/disaggregation/decode.py @@ -609,7 +609,11 @@ class DecodePreallocQueue: if not self.queue: return - if all(decode_req.waiting_for_input for decode_req in self.queue): + # Still poll if any receiver was aborted, otherwise it stays stuck. + if all(decode_req.waiting_for_input for decode_req in self.queue) and not any( + getattr(decode_req.kv_receiver, "conclude_state", None) == KVPoll.Failed + for decode_req in self.queue + ): return polls = poll_and_all_reduce( diff --git a/python/sglang/srt/managers/scheduler.py b/python/sglang/srt/managers/scheduler.py index 13a4946d6..67ac18d8c 100644 --- a/python/sglang/srt/managers/scheduler.py +++ b/python/sglang/srt/managers/scheduler.py @@ -3499,12 +3499,16 @@ class Scheduler( if recv_req.abort_all or decode_req.req.rid.startswith(recv_req.rid): logger.debug(f"Abort prealloc queue request. {decode_req.req.rid=}") decode_req.kv_receiver.abort() + if not isinstance(decode_req.req.finished_reason, FINISH_ABORT): + decode_req.req.finished_reason = FINISH_ABORT() # Abort requests waiting for kvcache to release tree cache for decode_req in self.disagg_decode_transfer_queue.queue: if recv_req.abort_all or decode_req.req.rid.startswith(recv_req.rid): logger.debug(f"Abort transfer queue request. {decode_req.req.rid=}") decode_req.kv_receiver.abort() + if not isinstance(decode_req.req.finished_reason, FINISH_ABORT): + decode_req.req.finished_reason = FINISH_ABORT() # Abort requests already retracted to CPU cache if self.disagg_decode_prealloc_queue.retracted_queue: diff --git a/python/sglang/srt/managers/tokenizer_manager.py b/python/sglang/srt/managers/tokenizer_manager.py index 83d8bf6b8..a22ea08c4 100644 --- a/python/sglang/srt/managers/tokenizer_manager.py +++ b/python/sglang/srt/managers/tokenizer_manager.py @@ -1491,7 +1491,15 @@ class TokenizerManager(TokenizerControlMixin, TokenizerManagerScoreMixin): pass def abort_request(self, rid: str = "", abort_all: bool = False): - if not abort_all and rid not in self.rid_to_state: + # Empty rid would startswith-match every request on the scheduler. + if not abort_all and not rid: + logger.warning("Ignore abort_request with empty rid and abort_all=False") + return + if ( + not abort_all + and self.server_args.tokenizer_worker_num == 1 + and rid not in self.rid_to_state + ): return req = AbortReq(rid=rid, abort_all=abort_all) self.send_to_scheduler.send_pyobj(req)