feat(ipc): control socket stage 2 - wake the render thread, brightness.set, plugin.reload (#720)

Control socket stage 2: the render thread wakes for queued commands (static screens ~1 ms, Vegas within one frame), brightness.set, and plugin.reload after a store update, with mailbox/restart fallbacks. Rig checks listed in the PR body are still to run.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Chuck
2026-10-01 20:37:54 -04:00
committed by GitHub
co-authored by Claude Opus 5.5
parent f72d69c2b0
commit 746dcfcadb
22 changed files with 1998 additions and 119 deletions
+116 -10
View File
@@ -8,6 +8,12 @@ where it reads the file mailbox (``DisplayController._poll_on_demand_requests``)
handing each command to the same code. Queries (``on_demand.status``) are
answered from a snapshot callable the display provides.
The queue also wakes the render thread: :meth:`ControlServer.wait_for_command`
is what it waits on in place of a sleep, so a command lands within a frame on
every kind of screen. An awaited command (``brightness.set``,
``plugin.reload``) carries a :class:`CommandOutcome` that the render thread
fills in; its connection thread waits for that, bounded, before answering.
Robustness rules, because this runs inside the display process:
* every connection has its own daemon thread, at most :data:`MAX_CLIENTS` at
@@ -39,10 +45,12 @@ import stat
import struct
import threading
import time
from dataclasses import dataclass
from typing import Any, Callable, Dict, FrozenSet, List, Mapping, Optional, Union
from dataclasses import dataclass, field
from typing import Any, Callable, Dict, FrozenSet, List, Mapping, Optional
from src.ipc.contract import (
AWAIT_SECONDS,
AWAITED_COMMANDS,
COMMANDS,
DEFAULT_SOCKET_DIR,
DEFAULT_SOCKET_PATH,
@@ -51,6 +59,7 @@ from src.ipc.contract import (
QUEUED_COMMANDS,
SUPPORTED_VERSIONS,
AckResult,
BrightnessSetArgs,
Command,
ErrorCode,
FrameReader,
@@ -58,7 +67,9 @@ from src.ipc.contract import (
HelloResult,
OnDemandStartArgs,
OnDemandStopArgs,
PluginReloadArgs,
ProtocolError,
QueuedArgs,
Request,
Response,
configured_socket_path,
@@ -100,19 +111,75 @@ _LISTEN_BACKLOG = 64
# -- queued work ---------------------------------------------------------------------
class CommandOutcome:
"""How an awaited command turned out, handed from the render thread back
to the connection thread that is waiting to answer.
The render thread calls :meth:`succeed` or :meth:`fail` once; the first
call wins. The connection thread may have stopped waiting already (it
answered ``pending``), and then nobody reads it.
"""
def __init__(self) -> None:
self._done = threading.Event()
self._lock = threading.Lock()
self.result: Optional[Dict[str, Any]] = None
self.error_code: Optional[str] = None
self.error_message = ''
@property
def done(self) -> bool:
return self._done.is_set()
def succeed(self, result: Mapping[str, Any]) -> None:
with self._lock:
if self._done.is_set():
return
self.result = dict(result)
self._done.set()
def fail(self, code: str, message: str) -> None:
with self._lock:
if self._done.is_set():
return
self.error_code = code
self.error_message = message
self._done.set()
def wait(self, timeout: float) -> bool:
return self._done.wait(timeout)
@dataclass(frozen=True)
class QueuedCommand:
"""A command waiting for the render thread."""
"""A command waiting for the render thread.
``outcome`` is set for an awaited command (``AWAITED_COMMANDS``): the
render thread reports through it, and the client's answer waits for it.
"""
request_id: str
cmd: str
args: Union[OnDemandStartArgs, OnDemandStopArgs]
args: QueuedArgs
received_at: float # time.time() when it was accepted
peer_uid: Optional[int] = None
outcome: Optional[CommandOutcome] = field(default=None, compare=False, repr=False)
def as_on_demand_request(self) -> Dict[str, Any]:
"""The mailbox-shaped payload the display's on-demand handler takes."""
if not isinstance(self.args, (OnDemandStartArgs, OnDemandStopArgs)):
raise TypeError(f'{self.cmd} is not an on-demand command')
return on_demand_request(self.request_id, self.args, self.received_at)
def succeed(self, result: Mapping[str, Any]) -> None:
"""Report success to a waiting client (a no-op for an acked command)."""
if self.outcome is not None:
self.outcome.succeed(result)
def fail(self, code: str, message: str) -> None:
"""Report failure to a waiting client (a no-op for an acked command)."""
if self.outcome is not None:
self.outcome.fail(code, message)
# -- peer credentials ------------------------------------------------------------------
@@ -247,8 +314,12 @@ class ControlServer:
max_clients: int = MAX_CLIENTS, io_timeout: float = IO_TIMEOUT_SECONDS,
message_timeout: float = MESSAGE_TIMEOUT_SECONDS,
idle_timeout: float = IDLE_TIMEOUT_SECONDS,
check_peer: bool = True):
check_peer: bool = True,
await_seconds: Optional[Mapping[str, float]] = None):
self.path = path
self._await_seconds: Dict[str, float] = dict(AWAIT_SECONDS)
if await_seconds:
self._await_seconds.update(await_seconds)
self._status_provider = status_provider
self._group = group
self._queue: 'queue.Queue[QueuedCommand]' = queue.Queue(maxsize=queue_size)
@@ -418,6 +489,17 @@ class ControlServer:
"""Cheap check for queued commands, for the render thread's fast path."""
return self._pending.is_set()
def wait_for_command(self, timeout: float) -> bool:
"""Block up to ``timeout`` seconds for a queued command; True if one is.
The render thread waits here instead of sleeping, in the dwell and
on a static screen, so a command wakes it at once. It is a timed
wait on an Event: no polling, and nothing more than the sleep it
replaces when no command comes. The flag stays set until drain(),
so a caller that does not drain would return at once every time.
"""
return self._pending.wait(timeout)
def drain(self) -> List[QueuedCommand]:
"""Every queued command, oldest first. Called from the render thread."""
commands: List[QueuedCommand] = []
@@ -596,11 +678,13 @@ class ControlServer:
v=request.v)
return Response.success(request.id, self._status_provider(), v=request.v)
if request.cmd in QUEUED_COMMANDS and isinstance(args, (OnDemandStartArgs,
OnDemandStopArgs)):
if request.cmd in QUEUED_COMMANDS and isinstance(args, (
OnDemandStartArgs, OnDemandStopArgs, BrightnessSetArgs, PluginReloadArgs)):
awaited = request.cmd in AWAITED_COMMANDS
command = QueuedCommand(request_id=request.id, cmd=request.cmd, args=args,
received_at=time.time(),
peer_uid=peer.uid if peer is not None else None)
peer_uid=peer.uid if peer is not None else None,
outcome=CommandOutcome() if awaited else None)
try:
self._queue.put_nowait(command)
except queue.Full:
@@ -610,15 +694,37 @@ class ControlServer:
'the display is not taking commands right now',
v=request.v)
self._pending.set()
logger.info("Control socket accepted %s %s", request.cmd, request.id)
if command.outcome is not None:
return self._await_outcome(request, command.outcome)
ack: AckResult = {'accepted': True, 'request_id': request.id,
'queued': self._queue.qsize()}
logger.info("Control socket accepted %s %s", request.cmd, request.id)
return Response.success(request.id, dict(ack), v=request.v)
# A command in COMMANDS with no handler here is a bug in this module.
return Response.failure(request.id, ErrorCode.INTERNAL,
f'{request.cmd} is not implemented', v=request.v)
def _await_outcome(self, request: Request, outcome: CommandOutcome) -> Response:
"""Answer an awaited command once the render thread has applied it.
Waits on this connection's thread, never the render thread's. If the
render thread does not get to it in time, the answer is ``pending``:
the command stays queued and is still applied, so a client treats
that as "not known to be done" rather than as a refusal.
"""
timeout = self._await_seconds.get(request.cmd, 0.0)
if not outcome.wait(timeout):
logger.warning("Control socket: %s %s not applied within %.1fs; answering pending",
request.cmd, request.id, timeout)
return Response.failure(request.id, ErrorCode.PENDING,
f'accepted, but not applied within {timeout:g}s; '
'the display will still apply it', v=request.v)
if outcome.error_code is not None:
return Response.failure(request.id, outcome.error_code, outcome.error_message,
v=request.v)
return Response.success(request.id, outcome.result or {}, v=request.v)
def start_control_server(status_provider: Optional[StatusProvider] = None,
cache_dir: Optional[str] = None,
@@ -637,7 +743,7 @@ def start_control_server(status_provider: Optional[StatusProvider] = None,
__all__ = [
'ControlServer', 'PeerCredentials', 'QueuedCommand', 'StatusProvider',
'CommandOutcome', 'ControlServer', 'PeerCredentials', 'QueuedCommand', 'StatusProvider',
'peer_allowed', 'peer_credentials', 'process_groups', 'resolve_socket_group',
'server_socket_path', 'start_control_server', 'PROTOCOL_VERSION',
]