@@ -142,6 +142,7 @@ def __init__(self, cfg: FDConfig, start_queue=True, use_async_llm=False):
142142
143143 self .is_paused = False # pause request generation
144144 self ._pause_cond = threading .Condition ()
145+ self ._rejecting_new_requests = False # blocks new requests during abort drain
145146
146147 self ._ctrl_output_queues = {}
147148 self ._ctrl_response_mailboxes = collections .defaultdict (collections .OrderedDict )
@@ -1325,7 +1326,7 @@ def _insert_zmq_task_to_scheduler(self):
13251326 trace_print (LoggingEventName .REQUEST_QUEUE_START , data ["request_id" ], data .get ("user" , "" ))
13261327 self .llm_logger .debug (f"Receive request from api server: { request } " )
13271328
1328- if self .is_paused :
1329+ if self .is_paused or self . _rejecting_new_requests :
13291330 self .llm_logger .warning (f"Engine is paused, drop request: { request } " )
13301331 self ._send_error_response (
13311332 request .request_id ,
@@ -1445,12 +1446,14 @@ def _control_pause(self, control_request: ControlRequest):
14451446 if self .is_paused :
14461447 self .llm_logger .info ("Engine is already paused, no need to pause again." )
14471448 return
1448- self .is_paused = True
14491449
1450- self .llm_logger .info ("Abort running requests." )
1450+ # Phase 1: Block new requests but keep scheduling loop running
1451+ # (scheduling loop must continue to process _trigger_abort)
1452+ self ._rejecting_new_requests = True
14511453
14521454 self .resource_manager .log_status ()
1453- # preempted all running reqs. preempted reqs will be append to ResourceManager.waiting queue
1455+
1456+ # Wait for current worker batch to complete
14541457 timeout , count = 60 , 0
14551458 while self .engine_worker_queue .exist_tasks ():
14561459 time .sleep (0.001 )
@@ -1461,21 +1464,39 @@ def _control_pause(self, control_request: ControlRequest):
14611464 error_msg = f"Emptying engine worker queue timed out after { timeout } seconds, worker may hanged!"
14621465 self .llm_logger .error (error_msg )
14631466 raise Exception (error_msg )
1464- running_reqs = self .resource_manager .preempted_all ()
1465- if len (running_reqs ) > 0 :
1466- self .llm_logger .info (f"Total { len (running_reqs )} requests need to be aborted." )
1467- self .resource_manager .get_real_bsz ()
1468- self .engine_worker_queue .put_tasks ((running_reqs , self .resource_manager .real_bsz ))
1469- self .resource_manager .wait_worker_inflight_requests_finish (timeout = 60 )
1470- # self.engine_worker_queue.clear_data()
1467+
1468+ # Phase 2: Trigger abort for ALL known requests
1469+ # Scheduling loop picks them up via _trigger_abort when they enter resource_manager
1470+ all_req_ids = list (
1471+ set (self .resource_manager .requests .keys ()) | set (self .scheduler .requests .keys ())
1472+ )
1473+ self .llm_logger .info (f"Pause: aborting { len (all_req_ids )} total requests." )
1474+ if all_req_ids :
1475+ self .resource_manager .add_abort_req_ids (all_req_ids )
1476+
1477+ # Phase 3: Wait for resource_manager to drain
1478+ # (all requests that enter rm will be caught by _trigger_abort)
1479+ self ._wait_inflight_drained (timeout = 30 )
1480+
1481+ # Phase 4: Fully pause the scheduling loop
1482+ with self ._pause_cond :
1483+ self .is_paused = True
1484+ self ._rejecting_new_requests = False
1485+
1486+ # Handle any requests remaining in scheduler (never pulled into rm)
1487+ remaining = set (self .scheduler .requests .keys ())
1488+ if remaining :
1489+ self .llm_logger .info (f"Pause: { len (remaining )} scheduler-only requests, sending abort directly." )
1490+ for req_id in remaining :
1491+ self ._send_error_response (req_id , "Aborted" , error_code = 200 )
1492+ self .resource_manager .waiting_abort_req_id_set .discard (req_id )
1493+
1494+ # Phase 5: Wait for output queue to be consumed (responses sent via ZMQ)
1495+ self ._wait_output_queue_empty (timeout = 5 )
1496+
1497+ # Phase 6: Reset
14711498 self .token_processor .clear_data ()
14721499 self .resource_manager .log_status ()
1473-
1474- # abort inflight requests to user
1475- inflight_requests = self .scheduler .get_inflight_requests ()
1476- self .llm_logger .info (f"Abort inflight requests (total { len (inflight_requests )} )." )
1477- for req in inflight_requests :
1478- self ._send_error_response (req .request_id , "Request is aborted since engine is paused." )
14791500 self .scheduler .reset ()
14801501
14811502 if envs .ENABLE_V1_KVCACHE_MANAGER :
@@ -1500,6 +1521,44 @@ def _control_pause(self, control_request: ControlRequest):
15001521 self .llm_logger .info ("Successfully paused request generation." )
15011522 return None
15021523
1524+ def _wait_inflight_drained (self , timeout = 30 ):
1525+ """
1526+ Wait until resource_manager.requests is completely empty.
1527+ All requests entering rm are caught by _trigger_abort in the scheduling loop.
1528+ """
1529+ start_time = time .time ()
1530+ while time .time () - start_time < timeout :
1531+ if not self .resource_manager .requests :
1532+ self .llm_logger .info ("All inflight requests drained." )
1533+ return
1534+ elapsed = time .time () - start_time
1535+ if elapsed > 10 :
1536+ self .llm_logger .warning (
1537+ f"Abort drain slow ({ elapsed :.1f} s): "
1538+ f"{ len (self .resource_manager .requests )} requests remaining"
1539+ )
1540+ time .sleep (0.005 )
1541+ self .llm_logger .warning (
1542+ f"_wait_inflight_drained timed out after { timeout } s, "
1543+ f"{ len (self .resource_manager .requests )} requests remaining"
1544+ )
1545+
1546+ def _wait_output_queue_empty (self , timeout = 5 ):
1547+ """
1548+ Wait until scheduler output queue is consumed by the output thread.
1549+ Ensures all abort responses have been sent via ZMQ before reset.
1550+ """
1551+ start_time = time .time ()
1552+ while time .time () - start_time < timeout :
1553+ with self .scheduler .mutex :
1554+ has_pending = bool (self .scheduler .responses ) or bool (
1555+ self .scheduler .batch_responses_per_step
1556+ )
1557+ if not has_pending :
1558+ return
1559+ time .sleep (0.005 )
1560+ self .llm_logger .warning (f"_wait_output_queue_empty timed out after { timeout } s" )
1561+
15031562 def _control_resume (self , control_request : ControlRequest ) -> Optional [dict ]:
15041563 """Control function for resuming request generation.
15051564
0 commit comments