diff --git a/python/sglang/srt/managers/scheduler.py b/python/sglang/srt/managers/scheduler.py index 68b6af232..4f6032fa6 100644 --- a/python/sglang/srt/managers/scheduler.py +++ b/python/sglang/srt/managers/scheduler.py @@ -137,6 +137,7 @@ from sglang.srt.managers.io_struct import ( ExpertDistributionReq, ExpertDistributionReqOutput, ExpertDistributionReqType, + FinishReasonDict, FlushCacheReqInput, FreezeGCReq, GetInternalStateReq, @@ -2821,13 +2822,13 @@ class Scheduler( and req.priority is not None and self.abort_on_priority_when_disabled ): - abort_req = AbortReq( + abort_req = _make_abort_req( + req, finished_reason={ "type": "abort", "status_code": HTTPStatus.SERVICE_UNAVAILABLE, "message": "Using priority is disabled for this server. Please send a new request without a priority.", }, - rid=req.rid, ) req.time_stats.trace_ctx.abort(abort_info=abort_req.finished_reason) self.ipc_channels.send_to_tokenizer.send_output(abort_req, req) @@ -2870,13 +2871,13 @@ class Scheduler( message = "The request is aborted by a higher priority request." self.ipc_channels.send_to_tokenizer.send_output( - AbortReq( + _make_abort_req( + req_to_abort, finished_reason={ "type": "abort", "status_code": HTTPStatus.SERVICE_UNAVAILABLE, "message": message, }, - rid=req_to_abort.rid, ), req_to_abort, ) @@ -2896,13 +2897,13 @@ class Scheduler( # Release prefetch events associated with the request self.tree_cache.release_aborted_request(req.rid) self.ipc_channels.send_to_tokenizer.send_output( - AbortReq( + _make_abort_req( + req, finished_reason={ "type": "abort", "status_code": HTTPStatus.SERVICE_UNAVAILABLE, "message": "Request waiting timeout reached.", }, - rid=req.rid, ), req, ) @@ -3035,7 +3036,7 @@ class Scheduler( self.chunked_req = None self._pending_chunked_abort_req = None - self.ipc_channels.send_to_tokenizer.send_output(AbortReq(rid=req.rid), req) + self.ipc_channels.send_to_tokenizer.send_output(_make_abort_req(req), req) logger.debug(f"Abort chunked prefill request. {req.rid=}") def _build_hisparse_decode_batch(self, reqs): @@ -3617,10 +3618,7 @@ class Scheduler( for req in reqs_to_abort: abort_reason: FINISH_ABORT = req.to_finish self.ipc_channels.send_to_tokenizer.send_output( - AbortReq( - finished_reason=abort_reason.to_json(), - rid=req.rid, - ), + _make_abort_req(req, finished_reason=abort_reason.to_json()), req, ) @@ -4604,7 +4602,7 @@ class Scheduler( if self.enable_hicache_storage: # to release prefetch events associated with the request self.tree_cache.release_aborted_request(req.rid) - self.ipc_channels.send_to_tokenizer.send_output(AbortReq(rid=req.rid), req) + self.ipc_channels.send_to_tokenizer.send_output(_make_abort_req(req), req) # For disaggregation decode mode, the request in the waiting queue has KV cache allocated. if self.disaggregation_mode == DisaggregationMode.DECODE: release_kv_cache(req, self.tree_cache) @@ -4637,7 +4635,7 @@ class Scheduler( if self.enable_hicache_storage: self.tree_cache.release_aborted_request(req.rid) self.ipc_channels.send_to_tokenizer.send_output( - AbortReq(rid=req.rid), req + _make_abort_req(req), req ) if ( req.req_pool_idx is not None @@ -4713,7 +4711,7 @@ class Scheduler( get_disagg().disaggregation_decode_retraction_backup, ) self.ipc_channels.send_to_tokenizer.send_output( - AbortReq(rid=decode_req.rid), decode_req + _make_abort_req(decode_req), decode_req ) else: remaining_retracted.append(decode_req) @@ -5235,3 +5233,9 @@ def run_scheduler_process( # and the synchronize() in destroy() could itself hang. if scheduler.gracefully_exit: scheduler.release_host_resources() + + +def _make_abort_req( + req: Req, finished_reason: Optional[FinishReasonDict] = None +) -> AbortReq: + return AbortReq(rid=req.rid, finished_reason=finished_reason)