02.10.01 — WebDriversPool[Driver]#
Fichier source :
src/ocarina/infra/drivers_pool.pyPool thread-safe de drivers, max concurrence garantie par sémaphore, warmup asynchrone surveillé, shutdown propre. Pas de réutilisation : chaque
acquire()détruit son driver à la fin (état propre).
Constructeur#
@final
class WebDriversPool[Driver]:
def __init__(
self,
create_driver: Thunk[BuiltWebDriver[Driver]],
max_size: int,
warmup_timeout: float | None = None,
) -> None:
self._create_driver = create_driver
self._pool: Queue[BuiltWebDriver[Driver]] = Queue(max_size)
self._semaphore = Semaphore(max_size)
self._warmup_timeout = (
warmup_timeout if warmup_timeout is not None and warmup_timeout > 0.1
else 60.0 * 5
)| Primitive | Rôle |
|---|---|
Queue(max_size) | File des drivers pré-créés disponibles. |
Semaphore(max_size) | Garantit que jamais plus de max_size drivers ne vivent simultanément (qu’ils soient en queue ou acquis). |
warmup_timeout | Délai max sans progression de warmup avant de lever WarmupTimeoutError. Par défaut 300s. |
acquire()#
@contextmanager
def acquire(self) -> Iterator[Driver]:
try:
driver, dispose = self._pool.get_nowait()
except Empty:
self._semaphore.acquire()
try:
driver, dispose = self._create_driver()
except Exception:
self._semaphore.release()
raise
try:
yield driver
finally:
with suppress(Exception):
dispose()
self._semaphore.release() acquire()
│
▼
┌──────────────────────────────────┐
│ pool.get_nowait() │── OK ─► driver, dispose ← cas warmup
└──────────────┬───────────────────┘
│ Empty
▼
┌──────────────────────────────────┐
│ sem.acquire() │ ← bloque si N drivers vivants
└──────────────┬───────────────────┘
▼
┌──────────────────────────────────┐
│ create_driver() │── leve ─► sem.release(); raise
└──────────────┬───────────────────┘
▼
┌──────────────────────────────────┐
│ yield driver │ ← caller utilise
└──────────────┬───────────────────┘
▼
finally :
with suppress : dispose() ← driver détruit
sem.release() ← rend une place1. Pas de réutilisation#
Queue.get_nowait() consomme l’entrée. Une fois sorti, le driver n’est jamais remis dans la queue. À la fin (finally), il est disposé.
Aucun driver ne sert pour 2 tests.
2. Pas de leak de sémaphore#
- Si
pool.get_nowait()retourne un driver : on l’utilise, on dispose, on release. ✅ - Si
pool.get_nowait()lèveEmpty: on acquire le sem, on crée, on utilise, on dispose, on release. ✅ - Si
create_driver()lève : on release le sem avant de propager. ✅ - Si
dispose()lève : on suppress + on release. ✅
Tous les chemins libèrent le sémaphore. Pas de fuite.
3. Concurrence garantie ≤ max_size#
Le sémaphore est acquis avant la création, donc tant que N drivers sont vivants (créés ou en queue), le N+1e appelant attend.
warmup()#
def warmup(self) -> None:
progress = {"count": 0}
lock = threading.Lock()
stop_event = threading.Event()
def worker() -> None:
while not self._pool.full() and not stop_event.is_set():
if not self._semaphore.acquire(blocking=False):
break
try:
driver, dispose = self._create_driver()
self._pool.put((driver, dispose))
with lock:
progress["count"] += 1
except Exception:
self._semaphore.release()
break
t = threading.Thread(target=worker, daemon=True)
t.start()
start_global = time.monotonic()
last_progress = 0
while t.is_alive():
time.sleep(0.5)
with lock:
current = progress["count"]
if current != last_progress:
last_progress = current
start_global = time.monotonic()
if time.monotonic() - start_global > self._warmup_timeout:
stop_event.set()
msg = (
"Warmup stalled (no progress detected)."
" Some browser processes may still be running."
" Please check your system (Activity Monitor / Task Manager / Dock)"
" and close any remaining browser instances."
)
self.shutdown()
raise WarmupTimeoutError(msg)
t.join()| Thread | Rôle |
|---|---|
| Worker (daemon) | Crée des drivers en boucle tant que la queue n’est pas pleine et que stop_event n’est pas set. À chaque succès, incrémente progress["count"]. |
| Watchdog (thread principal) | Boucle, échantillonne progress["count"] toutes les 0.5s. Si la valeur change, reset le timer. Si le timer dépasse warmup_timeout, stop_event.set() + raise WarmupTimeoutError. |
C’est un watchdog basé sur la progression, pas sur le temps absolu. Si la création d’un driver prend 30s mais que la progression avance, c’est OK. Si la création se bloque indéfiniment, la progression stagne → watchdog lève.
shutdown()#
def shutdown(self) -> None:
while not self._pool.empty():
_, dispose = self._pool.get_nowait()
with suppress(Exception):
dispose()
self._semaphore.release()- Dispose tous les drivers dans la queue (pas ceux acquis — ils sont gérés par leur
acquire()context). - Release les sémaphores correspondants.
suppress(Exception): si undispose()lève, on continue avec les suivants.
Appelé :
- Par
TestSuite.runen cas deWarmupTimeoutError. - Via
atexit.register(pool.shutdown)à la fin du process (cf.infra/selenium/create_drivers_pool.py).
WarmupTimeoutError#
class WarmupTimeoutError(Exception):
"""Raised when warmup is blocked (mostly some macOS edge cases)."""« macOS edge cases » : lorsqu’on désactive le mode headless et que l’on décide de fermer une fenêtre d’un navigateur piloté par Ocarina, macOS ne “tue” pas le process lié à la fenêtre.
Ce cas très spécifique a mis en avant la possibilité d’un hang infini dans le programme, et donc la mise en place d’un watchdog naïf.
Tests dédiés#
tests/scenarios/test_drivers_pool.py couvre :
- acquisition / release standard,
- warmup,
- détection WarmupTimeoutError quand le facteur stalle,
- comportement quand
create_driverlève (release du sem), - shutdown.