EngineCoreProc: wrapper de proceso en segundo plano con ZMQ
Responsabilidades
EngineCoreProc es una subclase de EngineCore (core.py:896-897), y ejecuta todo dentro de un proceso en segundo plano (background process) independiente. Su docstring lo deja claro: ZMQ-wrapper for running EngineCore in background process.. Hace tres cosas: durante el arranque del proceso hace handshake con el frontend para obtener las direcciones de los sockets ZMQ, levanta dos hilos de IO que puentean los sockets con dos queue.Queue (core.py:915-916), y luego run_busy_loop en el hilo principal va sacando requests de input_queue repetidamente → ejecuta step_fn() con scheduling y forward → mete los outputs en output_queue (core.py:1259-1267).
Esta capa es clave en la arquitectura multiproceso de v1. El proceso frontend (SyncMPClient o AsyncMPClient) envía vía ZMQ un EngineCoreRequest serializado; el hilo de input de EngineCoreProc lo recibe y hace put_nowait en input_queue, y el loop principal lo saca en un hueco del scheduling y lo despacha con _handle_client_request (core.py:1372-1405). Los outputs fluyen en sentido inverso: los EngineCoreOutputs que produce step_fn() se dejan en output_queue, el hilo de output los saca, los serializa y los envía de vuelta al frontend por ZMQ. El IO de los sockets ZMQ libera el GIL al enviar/recibir, así que esos dos hilos de IO pueden solaparse de verdad con el forward en GPU (core.py:974-1001); la serialización/deserialización también queda escondida detrás del forward.
Motivación de diseño
- Aislamiento de procesos:
EngineCoreProccorre en un proceso independiente; si el Python del frontend se cae no arrastra a los workers que ya están corriendo forward, y un OOM en un worker no se lleva por delante al proceso de API. Esa es la razón central por la que v1 saca EngineCore afuera: en v0LLMEnginehacía todo en el mismo proceso y un fallo en cualquier parte tumbaba todo. - Doble buffer con ZMQ + queue.Queue: el socket ZMQ no lo usa el loop principal directamente, hay una capa intermedia con
queue.Queue(core.py:915-916). El IO del socket vive en hilos de IO, el loop principal solo hacequeue.get(block=...), evitando que el loop principal se quede bloqueado en unrecvdel socket o arrastrado por un lock interno de ZMQ; al mismo tiempo, cuando el IO del socket libera el GIL, el loop principal puede retomarlo de inmediato para correr forward. - Dos hilos de IO:
process_input_socketsyprocess_output_sockets(core.py:980-1001) son independientes; el hilo de input se encarga de deserializar +put_nowait, el de output haceget+ serializar +send_multipart, y no se bloquean entre sí. - Handshake de arranque: en el constructor,
_perform_handshakes(core.py:932-938) obtiene unEngineZmqAddresses, y luego el hilo de input envía unEngineCoreReadyResponsepor ROUTER/DEALER para avisarle al frontend "ya estoy arriba, este es el max_model_len, num_gpu_blocks, kv_cache_size_tokens" (core.py:1525-1546). El frontend espera este mensaje ready enMPClient.__init__; si no llega, lo trata como fallo de arranque. run_engine_corecomo entrada del proceso: esta función es el target demultiprocessing.Process(target=run_engine_core, ...)(core.py:1154-1224). InstanciaEngineCoreProc, registra SIGTERM/SIGINT, y al final llamarun_busy_loop. Cualquier excepción cae enexcept Exception→_send_engine_dead→raise; cuando el proceso sale, enfinallyse llamaengine_core.shutdown()(core.py:1226-1242).
Archivos clave
EngineCoreProc class:896-910— definición declass EngineCoreProc(EngineCore):+ la constante de claseENGINE_CORE_DEAD = b"ENGINE_CORE_DEAD".EngineCoreProc.__init__:903-1009— levanta input_queue/output_queue, el receptor de tensor IPC,_perform_handshakes,super().__init__, y arranca los dos hilos de IO.queues + executor_fail_callback:915-919—input_queue = queue.Queue[...],output_queue = queue.Queue[...]; el callback de fallo del executor meteEXECUTOR_FAILEDen input_queue.IO threads startup:974-1001— arranca los dos daemon threadsprocess_input_socketsyprocess_output_sockets;ready_eventsolo se setea cuando el hilo de input queda ready.run_engine_core:1154-1224— entrada del proceso, eligeEngineCoreProcoDPEngineCoreProc, registra signal handlers, y llamarun_busy_loop.signal handlers:1204-1222—wakeup_enginemeteWAKEUPen input_queue; SIGTERM/SIGINT cambianshutdown_stateaREQUESTED.run_busy_loop:1259-1267—while self._handle_shutdown(): _process_input_queue(); _process_engine_step(), y al terminar el shutdown lanzaraise SystemExit._process_input_queue:1269-1298— sin trabajo bloquea eninput_queue.get(block=True), con trabajo hace un drain no bloqueante._process_engine_step:1300-1317— llamastep_fn(),output_queue.put_nowait(output),post_step; si no hubo forward pero el scheduler aún tiene trabajo, hacetime.sleep(0.001)para ceder el GIL._handle_client_request:1372-1405— despachaEngineCoreRequestType: WAKEUP/ADD/ABORT/UTILITY/EXECUTOR_FAILED._send_engine_dead:1470-1482— meteENGINE_CORE_DEADen output_queue y espera a que el hilo de output lo envíe antes de salir.process_input_sockets:1484-1587— cuerpo del hilo de input, usazmq.Pollersobre DEALER + un XSUB de coord opcional, deserializa y haceput_nowaiten input_queue.process_output_sockets:1589-1654— cuerpo del hilo de output, saca outputs conoutput_queue.get(), haceencoder.encode_into+send_multipart(copy=False, track=True)para enviarlos al frontend sin copiar.
Flujo de datos
Las requests entran por el socket de input de ZMQ, el hilo de input las deserializa y las mete en input_queue; el loop principal las saca en _process_input_queue y se las pasa a _handle_client_request para despacharlas:
# 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))Notar que ABORT se mete tanto en aborts_queue como en input_queue, porque abort es idempotente en el scheduler (core.py:1276-1279); así el loop principal puede también hacer drain de aborts_queue entre steps, evitando fugar requests ya abortados.
Cuando el loop principal recibe la request, _handle_client_request la despacha según 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 es un canal RPC genérico: el frontend envía el nombre del método + args msgspec, EngineCoreProc lo invoca por reflexión con getattr, y el resultado vuelve por output_queue. _invoke_utility_method también maneja el caso en el que se devuelve un Future (core.py:1440-1446): el output solo se envía cuando la utilidad async termina.
Límites y fallos
- Notificación de caída del proceso:
run_engine_coreen elexcept Exceptionllamaengine_core._send_engine_dead()(core.py:1229-1235), que mete enoutput_queuela cadena de bytesENGINE_CORE_DEAD(core.py:1470-1474); cuando el hilo de output la recibe, la difunde a todos los output sockets (core.py:1625-1628), y elMPClientdel frontend al recibirla lanza un error en vez de esperar indefinidamente. - Callback de fallo del executor: en
__init__se registra el closureexecutor_fail_callbackenmodel_executor(core.py:917-919); cuando un worker falla internamente, el callback mete unEXECUTOR_FAILEDeninput_queue, y el loop principal en_handle_client_requestlo recibe y lanzaraise RuntimeError("Executor failed.")(core.py:1400-1401); el proceso entero termina siendo capturado por el except derun_engine_corey cae en_send_engine_dead. - Modo shutdown:
_handle_shutdownmirashutdown_state(core.py:1324-1360); en estadoREQUESTED, sishutdown_timeout==0va por abort y mata de inmediato todas las requests en vuelo, si no es 0 va por drain y espera a que terminen todas. Durante el shutdown las nuevas ADD y UTILITY se rechazan explícitamente (core.py:1407-1432):_reject_add_in_shutdowndevuelve un output de abort al cliente, y_reject_utility_in_shutdowndevuelve unUtilityOutputconfailure_message="Server shutting down". - Muerte del hilo de input: en
__init__el loopready_event.wait(timeout=10)espera a que el hilo de input quede ready (core.py:1003-1009); si pasan 10 segundos sin ready y el hilo está muerto, lanzaraise RuntimeError("Input socket thread died during startup"); si solo no está ready pero sigue vivo, sigue esperando el mensaje READY del DP Coordinator. - Envío zero-copy + MessageTracker: el hilo de output usa
send_multipart(buffers, copy=False, track=True)(core.py:1646-1654), que devuelve unMessageTracker; los trackers que no terminaron de enviar, junto con referencias a outputs y buffers, se cuelgan en un dequepending, para que el GC no reclame el backing buffer que aún no se envió.reuse_buffersacota alen(sockets) + 1la cantidad de buffers reutilizables, evitando que la memoria crezca sin límite. - WAKEUP es no-op: el signal handler mete
WAKEUPen input_queue a través dewakeup_engine, solo para despertar al loop principal bloqueado eninput_queue.get(block=True)(core.py:1377-1378); no hace nada por sí mismo, el manejo real del shutdown lo hace_handle_shutdownmirandoshutdown_state. - finally restaura señales:
run_engine_coreenfinallyejecutasignal.signal(signal.SIGTERM, signal.SIG_DFL)(core.py:1236-1242), devuelve los signal handlers a su valor por defecto, y luego llamaengine_core.shutdown()para liberar model_executor / scheduler / estado distribuido.shutdowntambién llamagc.unfreeze()(core.py:652-656), para devolver al GC el startup heap que se había congelado congc.freeze()en el arranque; sin esto, el modo en el mismo proceso (unit test) filtra VRAM.
Resumen
EngineCoreProc es el muro de carga de la arquitectura multiproceso de v1: hereda de EngineCore, corre en un proceso independiente con ZMQ + dos queue.Queue + dos hilos de IO, y deja que el IO del socket y el forward en GPU se solapen de verdad; run_engine_core es la entrada del proceso, encargada de instanciar, registrar señales, atrapar excepciones y llamar shutdown. Toda la comunicación entre el frontend (SyncMPClient / AsyncMPClient) y el backend (EngineCore) pasa por esta capa. Para ver cómo se ensambla el núcleo del motor, continúa en /engine/engine-core; para ver la shell fina del frontend, en /engine/llm-engine; las tres implementaciones de cliente están en /client/inproc-mp, y el interior del scheduler en /scheduler/scheduler.