EngineCoreProc : enveloppe de processus arrière ZMQ
Responsabilités
EngineCoreProc est une sous-classe d'EngineCore (core.py:896-897) qui fait tourner tout son contenu dans un processus arrière séparé (background process). Le docstring est direct : ZMQ-wrapper for running EngineCore in background process.. Il fait trois choses : au démarrage, poignée de main avec le frontal pour récupérer l'adresse des sockets ZMQ ; lance deux threads d'IO pour relier socket et deux queue.Queue (core.py:915-916) ; puis run_busy_loop tourne sur le thread principal et boucle sur : prendre une requête depuis input_queue → step_fn() pour ordonnancer et lancer le forward → pousser les outputs dans output_queue (core.py:1259-1267).
Cette couche est la clé de l'architecture multi-processus v1. Le processus frontal (SyncMPClient ou AsyncMPClient) sérialise les EngineCoreRequest via ZMQ ; le thread input côté EngineCoreProc les reçoit et fait put_nowait dans input_queue, et la boucle principale, entre deux ordonnancements, les retire et les passe à _handle_client_request pour dispatch (core.py:1372-1405). Les outputs circulent en sens inverse : les EngineCoreOutputs produits par step_fn() sont poussés dans output_queue, le thread output les retire, les sérialise et les renvoie au frontal via ZMQ. Le socket IO ZMQ libère le GIL à l'envoi comme à la réception, donc ces deux threads d'IO se chevauchent vraiment avec le forward GPU (core.py:974-1001), et la sérialisation/désérialisation se trouve masquée derrière le forward.
Motivation de conception
- Isolation de processus :
EngineCoreProctourne dans un processus séparé, donc un crash Python frontal n'emporte pas les workers déjà engagés dans le GPU, et un OOM worker ne tue pas le processus d'API. C'est la raison centrale pour laquelle v1 extirpe EngineCore : en v0,LLMEnginefaisait tout dans le même processus, un crash à un endroit faisait tout tomber. - Double buffer ZMQ + queue.Queue : le socket ZMQ n'est pas utilisé directement par la boucle principale ; on intercale une
queue.Queue(core.py:915-916). Le socket IO vit dans le thread d'IO, la boucle principale ne fait quequeue.get(block=...), sans être bloquée par un recv socket ni par un verrou interne ZMQ ; et quand le socket IO libère le GIL, la boucle principale peut le récupérer pour lancer le forward. - Deux threads d'IO :
process_input_socketsetprocess_output_sockets(core.py:980-1001) sont indépendants : le thread input gère la désérialisation +put_nowait, le thread output gèreget+ sérialisation +send_multipart, sans se bloquer l'un l'autre. - Poignée de main au démarrage (handshake) : dans le constructeur,
_perform_handshakes(core.py:932-938) récupère unEngineZmqAddresses, puis le thread input envoie uneEngineCoreReadyResponsesur le ROUTER/DEALER pour signaler au frontal « je suis prêt, voilà max_model_len, num_gpu_blocks, kv_cache_size_tokens » (core.py:1525-1546). Côté frontal,MPClient.__init__attend ce message ready ; s'il n'arrive pas, le démarrage est considéré comme échoué. - Point d'entrée
run_engine_core: cette fonction est la cible demultiprocessing.Process(target=run_engine_core, ...)(core.py:1154-1224). Elle instancieEngineCoreProc, enregistre SIGTERM/SIGINT, puis appellerun_busy_loop. Toute exception passe parexcept Exception→_send_engine_dead→raise, et à la sortie lefinallyappelleengine_core.shutdown()(core.py:1226-1242).
Fichiers clés
EngineCoreProc class:896-910—class EngineCoreProc(EngineCore):+ constante de classeENGINE_CORE_DEAD = b"ENGINE_CORE_DEAD".EngineCoreProc.__init__:903-1009— crée input_queue/output_queue, tensor IPC receiver,_perform_handshakes,super().__init__, lance les deux threads d'IO.queues + executor_fail_callback:915-919—input_queue = queue.Queue[...],output_queue = queue.Queue[...], le callback d'échec d'executor pousseEXECUTOR_FAILEDdans input_queue.IO threads startup:974-1001— lancent les deux threads daemonprocess_input_socketsetprocess_output_sockets;ready_eventn'est set qu'après que le thread input est prêt.run_engine_core:1154-1224— point d'entrée du processus, choisitEngineCoreProcouDPEngineCoreProc, enregistre les handlers de signaux, appellerun_busy_loop.signal handlers:1204-1222—wakeup_enginepousseWAKEUPdans input_queue, SIGTERM/SIGINT passeshutdown_stateàREQUESTED.run_busy_loop:1259-1267—while self._handle_shutdown(): _process_input_queue(); _process_engine_step()puisraise SystemExità la sortie du shutdown._process_input_queue:1269-1298—input_queue.get(block=True)quand il n'y a rien à faire, drain non-bloquant sinon._process_engine_step:1300-1317— appellestep_fn(),output_queue.put_nowait(output),post_step; si aucun forward mais scheduler encore actif,time.sleep(0.001)rend le GIL._handle_client_request:1372-1405— dispatch parEngineCoreRequestType: WAKEUP/ADD/ABORT/UTILITY/EXECUTOR_FAILED._send_engine_dead:1470-1482— pousseENGINE_CORE_DEADdans output_queue et attend que le thread output l'envoie avant de sortir.process_input_sockets:1484-1587— corps du thread input, utilisezmq.Pollersur DEALER + XSUB coord optionnel, désérialise puisput_nowaitdans input_queue.process_output_sockets:1589-1654— corps du thread output, prend les outputs viaoutput_queue.get(),encoder.encode_into+send_multipart(copy=False, track=True)pour renvoyer au frontal en zero-copy.
Flux de données
La requête entre par le socket ZMQ input, est désérialisée par le thread input et poussée dans input_queue ; la boucle principale la retire dans _process_input_queue et la passe à _handle_client_request pour dispatch :
# vllm/v1/engine/core.py L1556-L1587
while True:
for input_socket, _ in poller.poll():
# (RequestType, RequestData)
type_frame, *data_frames = input_socket.recv_multipart(copy=False)
# NOTE(yongji): ignore READY message sent by DP coordinator
# that is used to notify newly started engines
if type_frame.buffer == b"READY":
assert input_socket == coord_socket
continue
request_type = EngineCoreRequestType(bytes(type_frame.buffer))
# Deserialize the request data.
request: Any
if request_type == EngineCoreRequestType.ADD:
req: EngineCoreRequest = add_request_decoder.decode(data_frames)
try:
request = self.preprocess_add_request(req)
except Exception:
self._handle_request_preproc_error(req)
continue
else:
request = generic_decoder.decode(data_frames)
if request_type == EngineCoreRequestType.ABORT:
# Aborts are added to *both* queues, allows us to eagerly
# process aborts while also ensuring ordering in the input
# queue to avoid leaking requests. This is ok because
# aborting in the scheduler is idempotent.
self.aborts_queue.put_nowait(request)
# Push to input queue for core busy loop.
self.input_queue.put_nowait((request_type, request))Notez que ABORT est poussé à la fois dans aborts_queue et input_queue, parce que l'abort est idempotent côté scheduler (core.py:1276-1279) ; cela permet à la boucle principale de drainer aborts_queue entre deux steps, pour éviter de fuite des requêtes déjà abortées.
Une fois la requête récupérée par la boucle principale, _handle_client_request dispatche selon EngineCoreRequestType (core.py:1372-1405) :
# vllm/v1/engine/core.py L1377-L1401
if request_type == EngineCoreRequestType.WAKEUP:
return
elif request_type == EngineCoreRequestType.ADD:
req, request_wave = request
if self._reject_add_in_shutdown(req):
return
self.add_request(req, request_wave)
elif request_type == EngineCoreRequestType.ABORT:
self.abort_requests(request)
elif request_type == EngineCoreRequestType.UTILITY:
client_idx, call_id, method_name, args = request
if self._reject_utility_in_shutdown(client_idx, call_id, method_name):
return
output = UtilityOutput(call_id)
# Lazily look-up utility method so that failure will be handled/returned.
get_result = lambda: (
(method := getattr(self, method_name))
and method(*self._convert_msgspec_args(method, args))
)
enqueue_output = lambda out: self.output_queue.put_nowait(
(client_idx, EngineCoreOutputs(utility_output=out))
)
self._invoke_utility_method(method_name, get_result, output, enqueue_output)
elif request_type == EngineCoreRequestType.EXECUTOR_FAILED:
raise RuntimeError("Executor failed.")UTILITY est un canal RPC générique : le frontal envoie le nom de méthode + args msgspec, EngineCoreProc l'invoque via getattr par réflexion, et le résultat repart par output_queue. _invoke_utility_method gère aussi le cas où la valeur renvoyée est un Future (core.py:1440-1446) ; l'output n'est renvoyé qu'une fois l'utilitaire asynchrone terminé.
Limites et échecs
- Notification de crash du processus :
run_engine_core, dans sonexcept Exception, appelleengine_core._send_engine_dead()(core.py:1229-1235) qui pousse la bytestringENGINE_CORE_DEADdansoutput_queue(core.py:1470-1474) ; le thread output, en voyant cette valeur spéciale, la diffuse à tous les sockets d'output (core.py:1625-1628) ; leMPClientfrontal lève alors une erreur au lieu d'attendre sans fin. - Callback d'échec d'Executor :
__init__enregistre la closureexecutor_fail_callbackauprès demodel_executor(core.py:917-919) ; en cas d'erreur côté worker, le callback pousseEXECUTOR_FAILEDdansinput_queue, et la boucle principale, dans_handle_client_request, lèveraise RuntimeError("Executor failed.")(core.py:1400-1401) ; le processus tout entier est ensuite attrapé par leexceptderun_engine_corequi déclenche_send_engine_dead. - Mode shutdown :
_handle_shutdownconsulteshutdown_state(core.py:1324-1360) ; à l'étatREQUESTED,shutdown_timeout==0déclenche le mode abort (tue immédiatement toutes les requêtes en vol), non nul déclenche le mode drain (attend la fin de toutes les requêtes). Les ADD et UTILITY qui arrivent pendant le shutdown sont explicitement rejetés (core.py:1407-1432) :_reject_add_in_shutdownrenvoie un output d'abort au client,_reject_utility_in_shutdownrenvoie unUtilityOutputavecfailure_message="Server shutting down". - Mort du thread input pendant le démarrage :
__init__boucle surready_event.wait(timeout=10)dans l'attente de la readiness du thread input (core.py:1003-1009) ; si au bout de 10 secondes le thread n'est toujours pas prêt et qu'il est mort, onraise RuntimeError("Input socket thread died during startup"); s'il est encore vivant mais pas prêt, on continue d'attendre le message READY du DP Coordinator. - Envoi zero-copy + MessageTracker : le thread output utilise
send_multipart(buffers, copy=False, track=True)(core.py:1646-1654) qui renvoie unMessageTracker; les trackers + références outputs + buffers non encore envoyés sont accrochés à une dequepending, pour empêcher le GC de récupérer les backing buffers avant leur envoi.reuse_bufferslimite le nombre de buffers réutilisés àlen(sockets) + 1pour éviter une croissance mémoire non bornée. - WAKEUP est un no-op : le handler de signaux pousse
WAKEUPdans input_queue viawakeup_engine, uniquement pour réveiller la boucle principale bloquée dansinput_queue.get(block=True)(core.py:1377-1378) ; en soi il ne fait rien, le vrai traitement du shutdown se fait dans_handle_shutdownen lisantshutdown_state. - Restauration des signaux dans le finally :
run_engine_core, dans sonfinally, faitsignal.signal(signal.SIGTERM, signal.SIG_DFL)(core.py:1236-1242) pour remettre le handler par défaut, puis appelleengine_core.shutdown()pour libérer model_executor / scheduler / état distribué.shutdownappelle aussigc.unfreeze()(core.py:652-656) pour rendre au GC le startup heap figé pargc.freeze()au démarrage, sinon le mode in-process (tests unitaires) fuit de la mémoire GPU.
Résumé
EngineCoreProc est le mur porteur de l'architecture multi-processus v1 : il hérite d'EngineCore, tourne dans un processus séparé avec ZMQ + deux queue.Queue + deux threads d'IO, de façon à ce que le socket IO et le forward GPU se chevauchent vraiment ; run_engine_core est le point d'entrée du processus, chargé de l'instantiation, de l'enregistrement des signaux, du filet de sécurité en cas d'exception et de shutdown. Toute communication entre fronts (SyncMPClient / AsyncMPClient) et back (EngineCore) passe par cette couche. Pour aller plus loin : /engine/engine-core pour voir comment le cœur du moteur lui-même est assemblé, /engine/llm-engine pour la coquille frontale, les trois implémentations client à /client/inproc-mp, et l'intérieur du scheduler à /scheduler/scheduler.
Voir la documentation officielle : Documentation vLLM · README