mirror of
https://github.com/ChuckBuilds/LEDMatrix.git
synced 2026-10-06 07:15:09 +00:00
CI caught a race in test_a_new_start_supersedes_the_pending_one: the dispatcher's worker could read the old start, the route cancel it and send the new one, and the worker's send of the old one land after it. cancel() now waits out a send already in flight (the worker holds a send lock while it reads the pending start and sends it), so whatever the caller sends next lands after it. Pinned by test_a_cancel_waits_for_a_send_in_flight, which fails without the wait. The Linux-only TestRealSocket test for a display that went away now expects the 202 and the start-timeout that follows, as the route answers since the background dispatcher. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
232 lines
9.7 KiB
Python
232 lines
9.7 KiB
Python
"""Deliver an on-demand start to a display that is not listening yet.
|
|
|
|
``POST /api/v3/display/on-demand/start`` can find no display on the control
|
|
socket: the service is stopped (the route starts it) or still loading its
|
|
plugins. The socket comes up only when the display's run loop starts, which
|
|
can take longer than a client waits -- the MQTT bridge gives up after 15 s.
|
|
So the route answers at once (``202``, ``status: "starting"``) and hands
|
|
the request to the one :class:`OnDemandDispatcher` of the web process, whose
|
|
worker thread sends it again until the display acknowledges it or
|
|
:data:`START_WAIT_SECONDS` pass.
|
|
|
|
* One request at a time: a newer start replaces the pending one, and a
|
|
stop cancels it (:meth:`OnDemandDispatcher.cancel`).
|
|
* Its outcome is :meth:`OnDemandDispatcher.status`, which
|
|
``GET /display/on-demand/status`` and ``/display/current-status`` report:
|
|
``starting`` while it waits, ``delivered`` once acknowledged (the
|
|
display's own state takes over from there), or ``error`` with
|
|
``start-timeout`` or the socket's reason.
|
|
|
|
Nothing is written to disk: the file mailbox that once carried such a
|
|
request is gone (docs/IPC_CONTROL_SOCKET.md, stage 5).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import threading
|
|
import time
|
|
from typing import Any, Callable, Dict, Optional
|
|
|
|
from src.ipc import client as control_client
|
|
from src.logging_config import get_logger
|
|
|
|
logger = get_logger(__name__)
|
|
|
|
#: How long the worker keeps sending a start before it gives up
|
|
#: (``start-timeout``). The socket comes up when the display's run loop
|
|
#: starts, after every plugin has loaded.
|
|
START_WAIT_SECONDS = 45.0
|
|
|
|
#: Gap between two sends while nothing is listening.
|
|
RETRY_INTERVAL = 0.5
|
|
|
|
#: How long a finished outcome (delivered, error, cancelled) is still
|
|
#: reported, so a client polling every few seconds sees it.
|
|
OUTCOME_SECONDS = 120.0
|
|
|
|
#: ``send(payload)`` hands the request to the display (the route's
|
|
#: ``_send_on_demand``) and raises ``ControlError`` when it does not take it.
|
|
Sender = Callable[[Dict[str, Any]], Any]
|
|
|
|
|
|
class OnDemandDispatcher:
|
|
"""One pending on-demand start, and the worker thread that delivers it."""
|
|
|
|
def __init__(self, send: Sender, *,
|
|
wait_seconds: Optional[float] = None,
|
|
retry_interval: Optional[float] = None,
|
|
clock: Callable[[], float] = time.monotonic,
|
|
wall_clock: Callable[[], float] = time.time):
|
|
self._send = send
|
|
self.wait_seconds = START_WAIT_SECONDS if wait_seconds is None else wait_seconds
|
|
self.retry_interval = RETRY_INTERVAL if retry_interval is None else retry_interval
|
|
self._clock = clock
|
|
self._wall = wall_clock
|
|
self._lock = threading.Lock()
|
|
# Held by the worker while it reads the pending start and sends it,
|
|
# so cancel() can wait out a send already in flight: a cancelled (or
|
|
# superseded) start must not reach the display after the request
|
|
# that replaced it. Never taken while holding _lock.
|
|
self._send_lock = threading.Lock()
|
|
self._wake = threading.Event()
|
|
# Bumped by every submit and cancel: a send that started under an
|
|
# older generation does not report its result as the current one.
|
|
self._generation = 0
|
|
self._pending: Optional[Dict[str, Any]] = None
|
|
self._deadline = 0.0
|
|
self._status: Optional[Dict[str, Any]] = None
|
|
self._finished_at: Optional[float] = None
|
|
self._thread: Optional[threading.Thread] = None
|
|
|
|
# -- the routes' side ------------------------------------------------------
|
|
|
|
def submit(self, payload: Dict[str, Any], wait_seconds: Optional[float] = None) -> None:
|
|
"""Deliver ``payload`` (an on-demand start) in the background,
|
|
replacing any start still pending."""
|
|
with self._lock:
|
|
self._generation += 1
|
|
superseded = self._pending
|
|
self._pending = dict(payload)
|
|
self._deadline = self._clock() + (self.wait_seconds if wait_seconds is None
|
|
else wait_seconds)
|
|
self._status = self._describe('starting', payload)
|
|
self._finished_at = None
|
|
if self._thread is None or not self._thread.is_alive():
|
|
self._thread = threading.Thread(target=self._run, name='on-demand-dispatch',
|
|
daemon=True)
|
|
self._thread.start()
|
|
if superseded is not None:
|
|
logger.info("On-demand start %s superseded by %s before the display took it",
|
|
superseded.get('request_id'), payload.get('request_id'))
|
|
self._wake.set()
|
|
|
|
def cancel(self, reason: str = 'cancelled') -> Optional[str]:
|
|
"""Drop the pending start (a stop arrived). Returns its request id,
|
|
or None when nothing was pending."""
|
|
with self._lock:
|
|
pending = self._pending
|
|
if pending is None:
|
|
return None
|
|
self._generation += 1
|
|
self._pending = None
|
|
self._finish(self._describe('idle', pending, last_event=reason))
|
|
self._wake.set()
|
|
# A send of it may be in flight (at most the client's timeout, 1 s):
|
|
# let it finish, so whatever the caller sends next lands after it.
|
|
with self._send_lock:
|
|
pass
|
|
logger.info("On-demand start %s cancelled before the display took it (%s)",
|
|
pending.get('request_id'), reason)
|
|
return pending.get('request_id')
|
|
|
|
def status(self) -> Optional[Dict[str, Any]]:
|
|
"""The pending start's state, or its outcome for OUTCOME_SECONDS
|
|
after it finished; None otherwise. In the shape of the display's
|
|
on-demand state (``active``, ``status``, ``error``, ...), plus
|
|
``source: "web"``."""
|
|
with self._lock:
|
|
if self._status is None:
|
|
return None
|
|
if (self._finished_at is not None
|
|
and self._clock() - self._finished_at > OUTCOME_SECONDS):
|
|
return None
|
|
return dict(self._status)
|
|
|
|
def pending(self) -> bool:
|
|
with self._lock:
|
|
return self._pending is not None
|
|
|
|
# -- the worker ------------------------------------------------------------
|
|
|
|
def _describe(self, status: str, payload: Dict[str, Any], error: Optional[str] = None,
|
|
last_event: Optional[str] = None) -> Dict[str, Any]:
|
|
return {
|
|
'active': False,
|
|
'status': status,
|
|
'error': error,
|
|
'last_event': last_event,
|
|
'request_id': payload.get('request_id'),
|
|
'plugin_id': payload.get('plugin_id'),
|
|
'mode': payload.get('mode'),
|
|
'duration': payload.get('duration'),
|
|
'pinned': bool(payload.get('pinned', False)),
|
|
'last_updated': self._wall(),
|
|
'source': 'web',
|
|
}
|
|
|
|
def _finish(self, status: Dict[str, Any]) -> None:
|
|
"""Record an outcome. Caller holds _lock."""
|
|
self._status = status
|
|
self._finished_at = self._clock()
|
|
|
|
def _run(self) -> None:
|
|
while True:
|
|
with self._send_lock:
|
|
with self._lock:
|
|
payload, generation = self._pending, self._generation
|
|
deadline = self._deadline
|
|
if payload is None:
|
|
self._thread = None
|
|
return
|
|
outcome, error = self._attempt(payload)
|
|
with self._lock:
|
|
if generation != self._generation:
|
|
continue # superseded or cancelled meanwhile
|
|
if outcome == 'retry' and self._clock() + self.retry_interval > deadline:
|
|
outcome, error = 'error', 'start-timeout'
|
|
if outcome == 'delivered':
|
|
self._pending = None
|
|
self._finish(self._describe('delivered', payload,
|
|
last_event='delivered'))
|
|
elif outcome == 'error':
|
|
self._pending = None
|
|
self._finish(self._describe('error', payload, error=error))
|
|
else:
|
|
self._wake.clear()
|
|
if outcome == 'delivered':
|
|
logger.info("On-demand start %s delivered once the display was listening",
|
|
payload.get('request_id'))
|
|
elif outcome == 'error':
|
|
logger.warning("On-demand start %s not delivered: %s",
|
|
payload.get('request_id'), error)
|
|
else:
|
|
# A submit or cancel wakes the wait at once.
|
|
self._wake.wait(self.retry_interval)
|
|
|
|
def _attempt(self, payload: Dict[str, Any]):
|
|
try:
|
|
self._send(payload)
|
|
except control_client.ControlError as e:
|
|
if control_client.display_not_listening(e):
|
|
return 'retry', None
|
|
return 'error', str(e.reason)
|
|
except Exception: # pylint: disable=broad-except
|
|
logger.exception("On-demand start %s: the control socket client failed",
|
|
payload.get('request_id'))
|
|
return 'error', 'internal'
|
|
return 'delivered', None
|
|
|
|
|
|
_dispatcher: Optional[OnDemandDispatcher] = None
|
|
_dispatcher_lock = threading.Lock()
|
|
|
|
|
|
def get_dispatcher(send: Sender) -> OnDemandDispatcher:
|
|
"""The web process's dispatcher, created on first use."""
|
|
global _dispatcher
|
|
with _dispatcher_lock:
|
|
if _dispatcher is None:
|
|
_dispatcher = OnDemandDispatcher(send)
|
|
return _dispatcher
|
|
|
|
|
|
def current() -> Optional[OnDemandDispatcher]:
|
|
"""The dispatcher if one was created, without creating one."""
|
|
return _dispatcher
|
|
|
|
|
|
def reset_for_tests() -> None:
|
|
global _dispatcher
|
|
with _dispatcher_lock:
|
|
_dispatcher = None
|