Skip to content

Commit e75ec66

Browse files
authored
fix(process-pool): bound shutdown latency (#469)
Fixes #297
1 parent 499a64f commit e75ec66

2 files changed

Lines changed: 118 additions & 4 deletions

File tree

openevolve/process_parallel.py

Lines changed: 59 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -332,6 +332,61 @@ def _run_iteration_worker(
332332
return SerializableResult(error=str(e), iteration=iteration)
333333

334334

335+
def _wait_for_processes(processes: tuple[mp.Process, ...], timeout: float) -> list[mp.Process]:
336+
"""Wait for process handles to observe worker exits without blocking indefinitely."""
337+
deadline = time.monotonic() + timeout
338+
alive = list(processes)
339+
while alive:
340+
next_alive = []
341+
for process in alive:
342+
try:
343+
process.join(timeout=0)
344+
if process.is_alive():
345+
next_alive.append(process)
346+
except (AssertionError, ValueError):
347+
continue
348+
alive = next_alive
349+
remaining = deadline - time.monotonic()
350+
if not alive or remaining <= 0:
351+
break
352+
time.sleep(min(0.001, remaining))
353+
return alive
354+
355+
356+
def _terminate_process_pool(executor: ProcessPoolExecutor) -> None:
357+
"""Cancel queued work and ensure all process-pool workers have exited."""
358+
# Python < 3.14 has no public force-shutdown API. Capture only this
359+
# executor's workers before shutdown clears its private process mapping.
360+
process_map = getattr(executor, "_processes", None) or {}
361+
processes = tuple(process_map.copy().values())
362+
terminate_workers = getattr(executor, "terminate_workers", None)
363+
364+
if callable(terminate_workers):
365+
terminate_workers()
366+
else:
367+
executor.shutdown(wait=False, cancel_futures=True)
368+
for process in processes:
369+
try:
370+
if process.is_alive():
371+
process.terminate()
372+
except (ProcessLookupError, ValueError):
373+
continue
374+
375+
surviving_processes = _wait_for_processes(processes, timeout=1.0)
376+
for process in surviving_processes:
377+
try:
378+
process.kill()
379+
except (ProcessLookupError, ValueError):
380+
continue
381+
382+
surviving_processes = _wait_for_processes(tuple(surviving_processes), timeout=1.0)
383+
if surviving_processes:
384+
logger.warning(
385+
"Process-pool workers did not exit: %s",
386+
[process.pid for process in surviving_processes],
387+
)
388+
389+
335390
class ProcessParallelController:
336391
"""Controller for process-based parallel evolution"""
337392

@@ -428,9 +483,10 @@ def stop(self) -> None:
428483
"""Stop the process pool"""
429484
self.shutdown_event.set()
430485

431-
if self.executor:
432-
self.executor.shutdown(wait=True)
433-
self.executor = None
486+
executor = self.executor
487+
self.executor = None
488+
if executor:
489+
_terminate_process_pool(executor)
434490

435491
logger.info("Stopped process pool")
436492

tests/test_process_parallel.py

Lines changed: 59 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,17 +4,26 @@
44

55
import asyncio
66
import os
7+
from pathlib import Path
78
import tempfile
89
import unittest
910
from unittest.mock import Mock, patch, MagicMock
1011
import time
11-
from concurrent.futures import Future
12+
from concurrent.futures import Future, ProcessPoolExecutor
13+
14+
15+
def _slow_test_worker(marker_path: str) -> str:
16+
Path(marker_path).write_text(str(os.getpid()))
17+
time.sleep(5)
18+
return "finished"
19+
1220

1321
# Set dummy API key for testing
1422
os.environ["OPENAI_API_KEY"] = "test"
1523

1624
from openevolve.config import Config, DatabaseConfig, EvaluatorConfig, LLMConfig, PromptConfig
1725
from openevolve.database import Program, ProgramDatabase
26+
from openevolve import process_parallel as process_parallel_module
1827
from openevolve.process_parallel import ProcessParallelController, SerializableResult
1928

2029

@@ -86,6 +95,55 @@ def test_controller_start_stop(self):
8695
self.assertIsNone(controller.executor)
8796
self.assertTrue(controller.shutdown_event.is_set())
8897

98+
def test_controller_stop_terminates_running_workers(self):
99+
"""Stopping the controller does not wait for stuck process-pool work."""
100+
controller = ProcessParallelController(self.config, self.eval_file, self.database)
101+
executor = ProcessPoolExecutor(max_workers=1)
102+
controller.executor = executor
103+
marker_path = os.path.join(self.test_dir, "worker.pid")
104+
future = executor.submit(_slow_test_worker, marker_path)
105+
106+
deadline = time.monotonic() + 5
107+
while not os.path.exists(marker_path) and time.monotonic() < deadline:
108+
time.sleep(0.01)
109+
self.assertTrue(os.path.exists(marker_path))
110+
worker_pid = int(Path(marker_path).read_text())
111+
112+
started = time.monotonic()
113+
controller.stop()
114+
elapsed = time.monotonic() - started
115+
116+
self.assertLess(elapsed, 1)
117+
self.assertIsNone(controller.executor)
118+
self.assertTrue(controller.shutdown_event.is_set())
119+
deadline = time.monotonic() + 1
120+
while not future.done() and time.monotonic() < deadline:
121+
time.sleep(0.01)
122+
self.assertTrue(future.done())
123+
with self.assertRaises(ProcessLookupError):
124+
os.kill(worker_pid, 0)
125+
126+
# Cleanup is idempotent after the executor reference is cleared.
127+
controller.stop()
128+
129+
def test_process_pool_shutdown_escalates_to_kill(self):
130+
"""Workers still alive after terminate are killed before returning."""
131+
process = Mock()
132+
process.is_alive.return_value = True
133+
executor = Mock(spec=["_processes", "shutdown"])
134+
executor._processes = {123: process}
135+
136+
with patch.object(
137+
process_parallel_module,
138+
"_wait_for_processes",
139+
side_effect=[[process], []],
140+
):
141+
process_parallel_module._terminate_process_pool(executor)
142+
143+
executor.shutdown.assert_called_once_with(wait=False, cancel_futures=True)
144+
process.terminate.assert_called_once_with()
145+
process.kill.assert_called_once_with()
146+
89147
def test_database_snapshot_creation(self):
90148
"""Test creating database snapshot for workers"""
91149
controller = ProcessParallelController(self.config, self.eval_file, self.database)

0 commit comments

Comments
 (0)