Source code for simvx.core.background
"""BackgroundSlot: run one slow call off the frame thread and poll for its result.
The frame loop is synchronous, so anything slow (a network request, a directory
walk, an expensive solve) freezes the game if it runs inline. A ``BackgroundSlot``
owns one worker thread and at most one in-flight call. The frame thread only ever
calls :meth:`submit` (returns immediately) and :meth:`poll` (non-blocking; hands
back a finished result or ``None``).
Policy:
- **Coalesce**: at most one call in flight. :meth:`submit` returns ``False`` and
does nothing while one is pending, so a caller that submits every frame still
issues one call at a time and never builds a backlog of stale work.
- **Graceful degrade**: a raising call is logged with its specific exception and
:meth:`poll` returns ``None``, so the caller keeps its last good value. There is
no bare ``except`` swallowing the cause.
- **Teardown**: :meth:`close` is idempotent and joins the worker. The worker is a
daemon thread, so a crashed game still exits.
::
slot: BackgroundSlot[Reply] = BackgroundSlot()
...
if not slot.pending:
slot.submit(lambda: expensive(snapshot))
result = slot.poll()
if result is not None:
self.last_good = result
"""
from __future__ import annotations
import logging
import threading
from collections.abc import Callable
from concurrent.futures import Future
log = logging.getLogger("simvx.core.background")
#: How long :meth:`BackgroundSlot.close` waits for a running call to return before
#: giving up on the join. The worker is a daemon, so a call that overruns this
#: never keeps the process alive.
_JOIN_TIMEOUT = 5.0
[docs]
class BackgroundSlot[T]:
"""Runs one callable off the frame thread at a time. Poll for the result."""
def __init__(self, *, thread_name: str = "simvx-background") -> None:
self._cond = threading.Condition()
self._queued: tuple[Callable[[], T], Future[T]] | None = None
self._inflight: Future[T] | None = None
self._closed = False
self._thread = threading.Thread(target=self._run, name=thread_name, daemon=True)
self._thread.start()
def _run(self) -> None:
"""Worker loop: take the one queued call, run it, resolve its future."""
while True:
with self._cond:
while self._queued is None and not self._closed:
self._cond.wait()
if self._queued is None:
return
fn, fut = self._queued
self._queued = None
if not fut.set_running_or_notify_cancel():
continue
try:
fut.set_result(fn())
except BaseException as exc: # noqa: BLE001 - reported through the future, logged in poll()
fut.set_exception(exc)
[docs]
@property
def pending(self) -> bool:
"""True while a submitted call is still running, so callers can gate cadence."""
with self._cond:
return self._inflight is not None and not self._inflight.done()
[docs]
def submit(self, fn: Callable[[], T]) -> bool:
"""Run ``fn`` on the worker thread, coalescing to one in-flight call.
Returns ``True`` if accepted, ``False`` if dropped because a call is already
in flight or the slot is closed. A dropped ``fn`` is never invoked.
"""
with self._cond:
if self._closed:
return False
if self._inflight is not None and not self._inflight.done():
return False
fut: Future[T] = Future()
self._inflight = fut
self._queued = (fn, fut)
self._cond.notify()
return True
[docs]
def poll(self) -> T | None:
"""Return a freshly finished result, or ``None`` if nothing finished or it failed.
Non-blocking: never waits on the worker. On success the slot is cleared and
the value handed back. On failure the exception is logged and ``None`` is
returned, so the caller retains its last good value. While a call is still
running, or nothing has been submitted, returns ``None``.
"""
with self._cond:
fut = self._inflight
if fut is None or not fut.done():
return None
self._inflight = None
if fut.cancelled():
return None
try:
return fut.result()
except Exception: # noqa: BLE001 - degrade gracefully, but log the specific cause
log.warning("BackgroundSlot call failed; keeping last good value", exc_info=True)
return None
[docs]
def cancel(self) -> None:
"""Drop the in-flight call without shutting the worker down.
A call the worker has already started runs to completion; its result is
discarded rather than handed to the next :meth:`poll`.
"""
with self._cond:
fut = self._inflight
self._inflight = None
self._queued = None
if fut is not None:
fut.cancel()
[docs]
def close(self) -> None:
"""Drop pending work, stop the worker and join it (idempotent)."""
with self._cond:
if self._closed:
return
self._closed = True
fut = self._inflight
self._inflight = None
self._queued = None
self._cond.notify_all()
if fut is not None:
fut.cancel()
self._thread.join(timeout=_JOIN_TIMEOUT)