[Fix] Wire the detokenizer soft watchdog into the multi-http-worker event loop (#31392)
This commit is contained in:
@@ -381,29 +381,31 @@ class MultiHttpWorkerDetokenizerMixin:
|
|||||||
def multi_http_worker_event_loop(self: DetokenizerManager):
|
def multi_http_worker_event_loop(self: DetokenizerManager):
|
||||||
"""The event loop that handles requests, for multi multi-http-worker mode"""
|
"""The event loop that handles requests, for multi multi-http-worker mode"""
|
||||||
self.socket_mapping = SocketMapping()
|
self.socket_mapping = SocketMapping()
|
||||||
|
# Watchdog wiring mirrors DetokenizerManager.event_loop: the watchdog is
|
||||||
|
# paused while waiting for input and fed once per processed message.
|
||||||
while True:
|
while True:
|
||||||
recv_obj = sock_recv(self.recv_from_scheduler)
|
with self.soft_watchdog.disable():
|
||||||
|
recv_obj = sock_recv(self.recv_from_scheduler)
|
||||||
output = self._request_dispatcher(recv_obj)
|
output = self._request_dispatcher(recv_obj)
|
||||||
if output is None:
|
if output is not None:
|
||||||
continue
|
# Fan out the output back to the originating tokenizer worker(s).
|
||||||
|
# In multi-detokenizer mode the upstream MultiDetokenizerRouter may
|
||||||
# Fan out the output back to the originating tokenizer worker(s).
|
# forward either batched or single requests, so handle both shapes.
|
||||||
# In multi-detokenizer mode the upstream MultiDetokenizerRouter may
|
if isinstance(recv_obj, BaseBatchReq):
|
||||||
# forward either batched or single requests, so handle both shapes.
|
for i, ipc_name in enumerate(recv_obj.http_worker_ipcs):
|
||||||
if isinstance(recv_obj, BaseBatchReq):
|
new_output = _handle_output_by_index(output, i)
|
||||||
for i, ipc_name in enumerate(recv_obj.http_worker_ipcs):
|
self.socket_mapping.send_output(
|
||||||
new_output = _handle_output_by_index(output, i)
|
ipc_name, new_output, is_tokenizer=True
|
||||||
|
)
|
||||||
|
elif isinstance(recv_obj, BaseReq):
|
||||||
self.socket_mapping.send_output(
|
self.socket_mapping.send_output(
|
||||||
ipc_name, new_output, is_tokenizer=True
|
recv_obj.http_worker_ipc, output, is_tokenizer=True
|
||||||
)
|
)
|
||||||
elif isinstance(recv_obj, BaseReq):
|
else:
|
||||||
self.socket_mapping.send_output(
|
raise ValueError(
|
||||||
recv_obj.http_worker_ipc, output, is_tokenizer=True
|
f"multi_http_worker_event_loop got unexpected req type {type(recv_obj)}"
|
||||||
)
|
)
|
||||||
else:
|
self.soft_watchdog.feed()
|
||||||
raise ValueError(
|
|
||||||
f"multi_http_worker_event_loop got unexpected req type {type(recv_obj)}"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
class MultiTokenizerRouter:
|
class MultiTokenizerRouter:
|
||||||
|
|||||||
Reference in New Issue
Block a user