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)