[Auto Sync] Update engine.py (20260102) (#16317)
Co-authored-by: github-actions[bot] <github-actions[bot]@users.noreply.github.com>
This commit is contained in:
@@ -819,6 +819,31 @@ def _set_envs_and_config(server_args: ServerArgs):
|
||||
mp.set_start_method("spawn", force=True)
|
||||
|
||||
|
||||
def _wait_for_scheduler_ready(
|
||||
scheduler_pipe_readers: List,
|
||||
scheduler_procs: List,
|
||||
) -> List[Dict]:
|
||||
"""Wait for the model to finish loading and return scheduler infos."""
|
||||
scheduler_infos = []
|
||||
for i in range(len(scheduler_pipe_readers)):
|
||||
try:
|
||||
data = scheduler_pipe_readers[i].recv()
|
||||
except EOFError:
|
||||
logger.error(
|
||||
f"Rank {i} scheduler is dead. Please check if there are relevant logs."
|
||||
)
|
||||
scheduler_procs[i].join()
|
||||
logger.error(f"Exit code: {scheduler_procs[i].exitcode}")
|
||||
raise
|
||||
|
||||
if data["status"] != "ready":
|
||||
raise RuntimeError(
|
||||
"Initialization failed. Please see the error messages above."
|
||||
)
|
||||
scheduler_infos.append(data)
|
||||
return scheduler_infos
|
||||
|
||||
|
||||
def _launch_scheduler_processes(
|
||||
server_args: ServerArgs,
|
||||
port_args: PortArgs,
|
||||
@@ -928,13 +953,13 @@ def _launch_subprocesses(
|
||||
# In multi-node cases, non-zero rank nodes do not need to run tokenizer or detokenizer,
|
||||
# so they can just wait here.
|
||||
|
||||
for reader in scheduler_pipe_readers:
|
||||
data = reader.recv()
|
||||
assert data["status"] == "ready"
|
||||
scheduler_infos = _wait_for_scheduler_ready(
|
||||
scheduler_pipe_readers, scheduler_procs
|
||||
)
|
||||
|
||||
if os.getenv("SGLANG_BLOCK_NONZERO_RANK_CHILDREN") == "0":
|
||||
# When using `Engine` as a Python API, we don't want to block here.
|
||||
return None, None, None, port_args
|
||||
return None, None, scheduler_infos, port_args
|
||||
|
||||
launch_dummy_health_check_server(
|
||||
server_args.host, server_args.port, server_args.enable_metrics
|
||||
@@ -945,7 +970,7 @@ def _launch_subprocesses(
|
||||
logger.error(
|
||||
f"Scheduler or DataParallelController {proc.pid} terminated with {proc.exitcode}"
|
||||
)
|
||||
return None, None, None, port_args
|
||||
return None, None, scheduler_infos, port_args
|
||||
|
||||
# Launch detokenizer process
|
||||
detoken_proc = mp.Process(
|
||||
@@ -968,23 +993,7 @@ def _launch_subprocesses(
|
||||
template_manager = None
|
||||
|
||||
# Wait for the model to finish loading
|
||||
scheduler_infos = []
|
||||
for i in range(len(scheduler_pipe_readers)):
|
||||
try:
|
||||
data = scheduler_pipe_readers[i].recv()
|
||||
except EOFError:
|
||||
logger.error(
|
||||
f"Rank {i} scheduler is dead. Please check if there are relevant logs."
|
||||
)
|
||||
scheduler_procs[i].join()
|
||||
logger.error(f"Exit code: {scheduler_procs[i].exitcode}")
|
||||
raise
|
||||
|
||||
if data["status"] != "ready":
|
||||
raise RuntimeError(
|
||||
"Initialization failed. Please see the error messages above."
|
||||
)
|
||||
scheduler_infos.append(data)
|
||||
scheduler_infos = _wait_for_scheduler_ready(scheduler_pipe_readers, scheduler_procs)
|
||||
|
||||
# Get back some info from scheduler to tokenizer_manager
|
||||
tokenizer_manager.max_req_input_len = scheduler_infos[0]["max_req_input_len"]
|
||||
|
||||
Reference in New Issue
Block a user