Base engine service client, responsible for managing EngineService lifecycle.
| 91 | |
| 92 | |
| 93 | class EngineServiceClient: |
| 94 | """ |
| 95 | Base engine service client, responsible for managing EngineService lifecycle. |
| 96 | """ |
| 97 | |
| 98 | def __init__(self, cfg, pid): |
| 99 | self.cfg = cfg |
| 100 | self.engine_process = None |
| 101 | self.engine_pid = pid |
| 102 | self._running = False |
| 103 | |
| 104 | llm_logger.info(f"EngineServiceClient initialized with engine_pid: {self.engine_pid}") |
| 105 | |
| 106 | async def start(self): |
| 107 | """Start engine service process""" |
| 108 | try: |
| 109 | # Start independent engine process |
| 110 | self._start_engine_process() |
| 111 | |
| 112 | # Wait for engine to be ready |
| 113 | if not self._wait_engine_ready(): |
| 114 | raise EngineError("Engine failed to start within timeout", error_code=500) |
| 115 | |
| 116 | self._running = True |
| 117 | llm_logger.info("EngineServiceClient started successfully") |
| 118 | |
| 119 | except Exception as e: |
| 120 | llm_logger.error(f"Failed to start EngineServiceClient: {e}") |
| 121 | raise |
| 122 | return True |
| 123 | |
| 124 | def _start_engine_process(self): |
| 125 | """Start engine process""" |
| 126 | try: |
| 127 | import multiprocessing |
| 128 | |
| 129 | self.shutdown_signal = multiprocessing.Value("i", 0) # 0=running, 1=shutdown |
| 130 | |
| 131 | def run_engine(): |
| 132 | engine = None |
| 133 | |
| 134 | def signal_handler(signum, frame): |
| 135 | llm_logger.info(f"Engine process received signal {signum}, initiating shutdown...") |
| 136 | if engine: |
| 137 | engine.running = False |
| 138 | |
| 139 | # Register signal handlers |
| 140 | signal.signal(signal.SIGTERM, signal_handler) |
| 141 | signal.signal(signal.SIGINT, signal_handler) |
| 142 | |
| 143 | try: |
| 144 | engine = EngineService(self.cfg, use_async_llm=True) |
| 145 | # Start engine with ZMQ service |
| 146 | engine.start(async_llm_pid=self.engine_pid) |
| 147 | |
| 148 | # Keep engine running until shutdown signal is received |
| 149 | while self.shutdown_signal.value == 0 and getattr(engine, "running", True): |
| 150 | time.sleep(0.5) |