fix nixl prefill crash make decode health check failed (#13657)
Signed-off-by: liluchang <liluchang@kingsoft.com>
This commit is contained in:
@@ -155,6 +155,9 @@ class NixlKVManager(CommonKVManager):
|
|||||||
self.max_failures = max(
|
self.max_failures = max(
|
||||||
get_int_env_var("SGLANG_DISAGGREGATION_HEARTBEAT_MAX_FAILURE", 2), 1
|
get_int_env_var("SGLANG_DISAGGREGATION_HEARTBEAT_MAX_FAILURE", 2), 1
|
||||||
)
|
)
|
||||||
|
self.waiting_timeout = get_int_env_var(
|
||||||
|
"SGLANG_DISAGGREGATION_WAITING_TIMEOUT", 300
|
||||||
|
)
|
||||||
self._start_heartbeat_checker_thread()
|
self._start_heartbeat_checker_thread()
|
||||||
else:
|
else:
|
||||||
raise ValueError(
|
raise ValueError(
|
||||||
@@ -186,16 +189,6 @@ class NixlKVManager(CommonKVManager):
|
|||||||
if response.status_code == 200:
|
if response.status_code == 200:
|
||||||
self.heartbeat_failures[bootstrap_addr] = 0
|
self.heartbeat_failures[bootstrap_addr] = 0
|
||||||
|
|
||||||
current_rooms = self.addr_to_rooms_tracker[
|
|
||||||
bootstrap_addr
|
|
||||||
].copy()
|
|
||||||
|
|
||||||
for bootstrap_room in current_rooms:
|
|
||||||
# Remove successful transfers from the tracker
|
|
||||||
if bootstrap_room not in self.transfer_statuses:
|
|
||||||
self.addr_to_rooms_tracker[bootstrap_addr].discard(
|
|
||||||
bootstrap_room
|
|
||||||
)
|
|
||||||
else:
|
else:
|
||||||
logger.info(
|
logger.info(
|
||||||
f"Attempting to reconnect to {bootstrap_addr}..."
|
f"Attempting to reconnect to {bootstrap_addr}..."
|
||||||
@@ -259,6 +252,9 @@ class NixlKVManager(CommonKVManager):
|
|||||||
f"Lost connection with prefill instance (bootstrap_addr: {failed_bootstrap_addr}), "
|
f"Lost connection with prefill instance (bootstrap_addr: {failed_bootstrap_addr}), "
|
||||||
f"{len(affected_rooms)} transfers affected"
|
f"{len(affected_rooms)} transfers affected"
|
||||||
)
|
)
|
||||||
|
for room in possible_affected_rooms:
|
||||||
|
logger.error(f"Let room {room} be failed due to prefill down")
|
||||||
|
self.update_status(room, KVPoll.Failed)
|
||||||
|
|
||||||
def check_status(self, bootstrap_room: int):
|
def check_status(self, bootstrap_room: int):
|
||||||
return self.request_status[bootstrap_room]
|
return self.request_status[bootstrap_room]
|
||||||
@@ -267,10 +263,13 @@ class NixlKVManager(CommonKVManager):
|
|||||||
if bootstrap_room not in self.request_status:
|
if bootstrap_room not in self.request_status:
|
||||||
self.request_status[bootstrap_room] = status
|
self.request_status[bootstrap_room] = status
|
||||||
else:
|
else:
|
||||||
# NOTE: The prefill engine could recv bootstrapping first
|
# NOTE: status is only allowed to be incremented unless it is KVPoll.Failed
|
||||||
self.request_status[bootstrap_room] = max(
|
if status == KVPoll.Failed:
|
||||||
self.request_status[bootstrap_room], status
|
self.request_status[bootstrap_room] = KVPoll.Failed
|
||||||
)
|
else:
|
||||||
|
self.request_status[bootstrap_room] = max(
|
||||||
|
self.request_status[bootstrap_room], status
|
||||||
|
)
|
||||||
|
|
||||||
def record_failure(self, bootstrap_room: int, failure_reason: str):
|
def record_failure(self, bootstrap_room: int, failure_reason: str):
|
||||||
pass
|
pass
|
||||||
@@ -752,6 +751,7 @@ class NixlKVReceiver(CommonKVReceiver):
|
|||||||
self.kv_mgr.addr_to_rooms_tracker[self.bootstrap_addr].add(
|
self.kv_mgr.addr_to_rooms_tracker[self.bootstrap_addr].add(
|
||||||
self.bootstrap_room
|
self.bootstrap_room
|
||||||
)
|
)
|
||||||
|
self.init_time = None
|
||||||
|
|
||||||
def init(
|
def init(
|
||||||
self,
|
self,
|
||||||
@@ -790,15 +790,35 @@ class NixlKVReceiver(CommonKVReceiver):
|
|||||||
)
|
)
|
||||||
|
|
||||||
self.started_transfer = True
|
self.started_transfer = True
|
||||||
|
self.init_time = time.time()
|
||||||
|
|
||||||
def poll(self) -> KVPoll:
|
def poll(self) -> KVPoll:
|
||||||
if self.conclude_state is not None:
|
if self.conclude_state is not None:
|
||||||
return self.conclude_state
|
return self.conclude_state
|
||||||
|
status = self.kv_mgr.check_status(self.bootstrap_room)
|
||||||
|
if status in (KVPoll.Success, KVPoll.Failed):
|
||||||
|
self.conclude_state = status
|
||||||
|
return status
|
||||||
if not self.started_transfer:
|
if not self.started_transfer:
|
||||||
return KVPoll.WaitingForInput # type: ignore
|
return KVPoll.WaitingForInput # type: ignore
|
||||||
|
|
||||||
|
now = time.time()
|
||||||
|
elapsed = now - self.init_time
|
||||||
|
|
||||||
|
if elapsed >= self.kv_mgr.waiting_timeout:
|
||||||
|
logger.error(f"Request {self.bootstrap_room} waiting_timeout")
|
||||||
|
self.kv_mgr.record_failure(
|
||||||
|
self.bootstrap_room,
|
||||||
|
f"Request {self.bootstrap_room} timed out after {elapsed:.1f}s in KVPoll.WaitingForInput",
|
||||||
|
)
|
||||||
|
self.conclude_state = KVPoll.Failed
|
||||||
|
return KVPoll.Failed
|
||||||
|
|
||||||
self.kv_mgr.update_transfer_status()
|
self.kv_mgr.update_transfer_status()
|
||||||
if self.kv_mgr.check_transfer_done(self.bootstrap_room): # type: ignore
|
if self.kv_mgr.check_transfer_done(self.bootstrap_room): # type: ignore
|
||||||
|
self.kv_mgr.addr_to_rooms_tracker[self.bootstrap_addr].discard(
|
||||||
|
self.bootstrap_room
|
||||||
|
)
|
||||||
# Check if the transfer failed
|
# Check if the transfer failed
|
||||||
if self.kv_mgr.transfer_statuses[self.bootstrap_room].is_failed():
|
if self.kv_mgr.transfer_statuses[self.bootstrap_room].is_failed():
|
||||||
self.conclude_state = KVPoll.Failed
|
self.conclude_state = KVPoll.Failed
|
||||||
|
|||||||
Reference in New Issue
Block a user