Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,36 @@ accepts both, but the store flags the old spelling as deprecated

## Unreleased

### Control socket stage 3: the display's state over the socket

- Two new commands, still protocol version 1. `state.get` returns a
versioned snapshot of what the display is doing: the current mode and
plugin, the on-demand session, the brightness, the plugin runtime
snapshot and the render loop's heartbeat age. With `since`/`epoch` it
returns a short "unchanged" answer. `state.subscribe` returns the same
snapshot, then pushes a `state` event on every change (always the latest
version) and a `tick` at least every 5 s. The display serves all of it
from memory (`StateHub` in `src/ipc/server.py`), and publishing never
waits for a reader. Subscribers have their own bound (4), separate from
the 8 request slots, and one that stops reading is dropped after the 2 s
IO timeout. See `docs/IPC_CONTROL_SOCKET.md`, "The state stream".
- The web interface holds one subscription per process
(`web_interface/display_state.py`). `/display/current-status`,
`/display/on-demand/status`, the plugin runtime fields of
`/plugins/installed` and `/plugins/state`, the reconciliations and
`/health`'s `display_loop` read it first. When the socket is missing (a
stopped or older display, Windows), they fall back to the cache keys and
the heartbeat file. Each answer has a `source` (`socket`, `cache` or
`heartbeat_file`). The stale and stalled rules from #726 apply the same
way to both.
- Fewer SD-card writes while the socket serves those readers.
`display_current_state` is written once a minute and on a flag change,
not on every mode change. The `plugin_runtime_snapshot` refresh goes from
60 s to 120 s. For a rotation of 15 s screens, that is 1.5 cache writes a
minute instead of 5. Both keys keep being written for one release.
- `RenderWatchdog.liveness()` reports the heartbeat age from memory.
`PluginRuntimeView` has a `source`, and `describe()` includes it.

### Web UI: four more tabs are ES-module pages (stage 2)

- Rotation, Operation History, Config Editor and Backup & Restore follow the
Expand Down
266 changes: 242 additions & 24 deletions docs/IPC_CONTROL_SOCKET.md

Large diffs are not rendered by default.

31 changes: 25 additions & 6 deletions docs/REST_API_REFERENCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -331,12 +331,19 @@ by the display process (stale after 120 seconds).
"data": {
"mode": "nfl_live",
"plugin_id": "football-scoreboard",
"last_updated": 1234567890.123
"last_updated": 1234567890.123,
"source": "socket"
}
}
```

When nothing has been published, every field is `null`.
When nothing has been published, every field is `null`. `source` is
`socket` when the answer came from the display's state stream over the
control socket ([IPC_CONTROL_SOCKET.md](IPC_CONTROL_SOCKET.md)), and `cache`
when it came from the `display_current_state` cache key (no socket: the
display is stopped or older, or this is Windows). A display whose render
loop has not refreshed its state for 120 seconds is reported with every
field `null`, either way.

### List Display Modes

Expand Down Expand Up @@ -413,11 +420,16 @@ Get the current on-demand display state.
"returncode": 0,
"stdout": "active",
"stderr": ""
}
},
"source": "socket"
}
}
```

`source` is `socket` (the display's state stream, with `remaining` worked
out at the time of the request) or `cache` (the `display_on_demand_state`
cache key).

With no on-demand request, `state` is
`{"active": false, "status": "idle", "last_updated": null}`.

Expand Down Expand Up @@ -552,7 +564,8 @@ List all installed plugins with their status and metadata.
"published_at": 1790000030.0,
"age_seconds": 12.4,
"stale_after": 180.0,
"heartbeat_age_seconds": 2.1
"heartbeat_age_seconds": 2.1,
"source": "socket"
}
}
}
Expand Down Expand Up @@ -581,7 +594,10 @@ the display is hung or died), `stopped` (the display shut down) or `unknown`
(nothing published yet). Unless it is `live`, every one of those fields is
`null`. `heartbeat_age_seconds` is the heartbeat's age when it was taken into
account, `null` otherwise (no heartbeat, as on the dev server, or one from
another process). Health and metrics are at [`/plugins/health`](#get-plugin-health)
another process). `runtime.source` is `socket` when the snapshot and the
heartbeat age came from the display's state stream over the control socket,
and `cache` when they came from the `plugin_runtime_snapshot` cache key and
the heartbeat file; the rules above are the same for both. Health and metrics are at [`/plugins/health`](#get-plugin-health)
and `/plugins/metrics`.

`vegas_participation` is what Vegas mode does with the plugin: `"scroll"`,
Expand Down Expand Up @@ -2265,7 +2281,10 @@ display snapshot. `data.status` is `healthy` or `degraded`, with
(with `heartbeat_age_seconds`), `stalled` (no heartbeat for 60s: the panel is
frozen even if the service is active; the status turns `degraded`), or
`not_reported` when the display writes none (not started yet, the dev server,
Windows), which does not affect the status.
Windows), which does not affect the status. Its `source` is `socket` when the
age came from the display's state stream over the control socket (measured
in memory by the display) and `heartbeat_file` when it came from
`/run/ledmatrix/display-heartbeat.json`.

Open even when the web login is on, for uptime monitors; a caller that is not
logged in (and has no token) then gets only `{"status": "success", "data":
Expand Down
114 changes: 100 additions & 14 deletions src/display_controller.py
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@
PluginReloadArgs,
PluginReloadResult,
)
from src.ipc.server import ControlServer, QueuedCommand, start_control_server
from src.ipc.server import ControlServer, QueuedCommand, StateHub, start_control_server
from src.vegas_mode.render_pipeline import SYNC_SEND_INTERVAL

# Get logger with consistent configuration
Expand All @@ -65,6 +65,13 @@
# treats display_current_state older than 120 s as unknown.
CURRENT_STATE_REFRESH_SECONDS = 30

# While the control socket serves the web interface's state readers
# (StateHub.readers_active), display_current_state is only their fallback:
# it is then rewritten at this interval and on a change of the flags, not on
# every mode change. Below the readers' 120 s max_age, so the fallback copy
# never reads as unknown.
CURRENT_STATE_RELAXED_REFRESH_SECONDS = 60

# How long startup will wait for plugins to fetch their first data before
# showing anything. Each plugin's update blocks for up to the executor's 30s
# timeout and they run one after another, so the uncapped total is the sum of
Expand Down Expand Up @@ -1430,18 +1437,60 @@ def _get_on_demand_remaining(self) -> Optional[float]:
return None
return max(0.0, expires_at - time.time())

#: The control socket's state stream (src/ipc/server.StateHub), while the
#: socket is served. Class-level default for controllers built without
#: __init__ (tests) and for a display with no socket.
_state_hub: Optional[StateHub] = None

def _current_mode_state(self) -> Dict[str, Any]:
"""What display_current_state and the socket's ``display`` section hold."""
return {
'mode': self.current_display_mode,
'plugin_id': self.mode_to_plugin_id.get(self.current_display_mode),
'mode_index': self.current_mode_index,
'total_modes': len(self.available_modes),
'on_demand_active': self.on_demand_active,
'is_display_active': self.is_display_active,
'last_updated': time.time(),
}

def _push_live_state(self, display_state: Optional[Dict[str, Any]] = None) -> None:
"""Hand the current mode and the brightness to the control socket's
state stream. In memory, no disk: the hub only bumps its version (and
wakes subscribers) when something other than ``last_updated`` changed.

Called on every pass of the publish points below, so
``display.last_updated`` doubles as the render thread's proof of life
for the socket's readers, as the cache key's max_age does today.
"""
hub = self._state_hub
if hub is None:
return
try:
hub.publish('display', display_state or self._current_mode_state(),
volatile=('last_updated',))
hub.publish('brightness', {
'brightness': getattr(self, '_normal_brightness', None),
'panel_brightness': getattr(self, 'current_brightness', None),
'dimmed': bool(getattr(self, 'is_dimmed', False)),
})
except Exception as err: # pylint: disable=broad-except
logger.debug("Could not publish the display state to the control socket: %s",
err, exc_info=True)

def _state_readers_on_socket(self) -> bool:
"""Is the control socket serving the web interface's state readers?"""
hub = self._state_hub
try:
return bool(hub is not None and hub.readers_active())
except Exception: # pylint: disable=broad-except
return False

def _publish_current_mode_state(self) -> None:
"""Publish the currently active display mode/plugin to cache for the web UI."""
try:
state = {
'mode': self.current_display_mode,
'plugin_id': self.mode_to_plugin_id.get(self.current_display_mode),
'mode_index': self.current_mode_index,
'total_modes': len(self.available_modes),
'on_demand_active': self.on_demand_active,
'is_display_active': self.is_display_active,
'last_updated': time.time(),
}
state = self._current_mode_state()
self._push_live_state(state)
self.cache_manager.set('display_current_state', state)
self._last_published_mode = self.current_display_mode
self._last_published_flags = self._current_state_flags()
Expand All @@ -1467,11 +1516,23 @@ def _publish_current_mode_state_if_changed(self) -> None:
priority, a single enabled plugin -- has to be republished or the UI
reports it as unknown. Otherwise this writes only on a change, not on
every render tick.

While the control socket serves the web interface's state readers,
the socket's in-memory copy is updated on every call and the cache
key is only their fallback: a mode change alone is then written at
the relaxed refresh (CURRENT_STATE_RELAXED_REFRESH_SECONDS), still
inside the readers' max_age. The flags are still written at once.
When the socket stops serving them, the next call writes a changed
mode again.
"""
if (self.current_display_mode != self._last_published_mode
relaxed = self._state_readers_on_socket()
refresh = CURRENT_STATE_RELAXED_REFRESH_SECONDS if relaxed else CURRENT_STATE_REFRESH_SECONDS
if ((not relaxed and self.current_display_mode != self._last_published_mode)
or self._current_state_flags() != getattr(self, '_last_published_flags', None)
or time.monotonic() - self._last_published_at >= CURRENT_STATE_REFRESH_SECONDS):
self._publish_current_mode_state()
or time.monotonic() - self._last_published_at >= refresh):
self._publish_current_mode_state() # pushes to the socket as well
else:
self._push_live_state()

def _on_demand_state(self) -> Dict[str, Any]:
"""The on-demand state as published to the cache and the control socket."""
Expand All @@ -1494,6 +1555,11 @@ def _publish_on_demand_state(self) -> None:
"""Publish current on-demand state to cache for external consumers."""
try:
state = self._on_demand_state()
hub = self._state_hub
if hub is not None:
# In memory, first: a subscriber hears the outcome of an
# on-demand command even if the cache write below fails.
hub.publish('on_demand', state, volatile=('last_updated', 'remaining'))
self.cache_manager.set('display_on_demand_state', state)
except (OSError, RuntimeError, ValueError, TypeError) as err:
logger.error("Failed to publish on-demand state: %s", err, exc_info=True)
Expand Down Expand Up @@ -1738,12 +1804,31 @@ def _start_control_server(self) -> None:
"""
if self._control_server is not None:
return
hub = StateHub(loop_probe=display_watchdog.watchdog.liveness)
try:
self._control_server = start_control_server(
status_provider=self._control_status,
cache_dir=getattr(self.cache_manager, 'cache_dir', None))
cache_dir=getattr(self.cache_manager, 'cache_dir', None),
state_hub=hub)
except Exception: # pylint: disable=broad-except
logger.exception("Control socket not started; using the file mailbox only")
if self._control_server is not None:
self._start_state_stream(hub)

def _start_state_stream(self, hub: StateHub) -> None:
"""Start publishing to the socket's state stream (``state.get`` and
``state.subscribe``): everything a reader would see, now, then on
every publish. Never raises; without it readers use the cache keys."""
try:
self._state_hub = hub
self._push_live_state()
hub.publish('on_demand', self._on_demand_state(),
volatile=('last_updated', 'remaining'))
publisher = getattr(self, '_plugin_runtime_publisher', None)
if publisher is not None:
publisher.attach_hub(hub)
except Exception: # pylint: disable=broad-except
logger.exception("Control socket state stream not started; readers use the cache")

def _control_status(self) -> Dict[str, Any]:
"""The socket's on_demand.status answer. Runs on the socket's thread: reads only."""
Expand Down Expand Up @@ -4532,6 +4617,7 @@ def cleanup(self):
except Exception as e:
logger.warning("Error closing the control socket: %s", e)
self._control_server = None
self._state_hub = None
# Stop the async update worker first so no in-flight update() call
# is still touching display/cache-backed resources while they're
# torn down below.
Expand Down
17 changes: 17 additions & 0 deletions src/display_watchdog.py
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,23 @@ def __init__(self, environ: Optional[Mapping[str, str]] = None,
def armed(self) -> bool:
return self._armed

def liveness(self) -> Dict[str, Any]:
"""The heartbeat, in memory: what the control socket's state stream
reports as ``loop``.

``heartbeat_age_seconds`` is the age of the render thread's last beat,
the beat that writes the heartbeat file, so it ages at the same rate
and is judged by the same ``HEARTBEAT_STALE_SECONDS``. None until the
loop has drawn its first frame. Any thread may call this: it only
reads two attributes.
"""
last = self._last_beat
age = None
if self._armed and last is not None:
age = max(self._clock() - last, 0.0)
return {'heartbeat_age_seconds': age, 'armed': self._armed,
'stale_after': HEARTBEAT_STALE_SECONDS}

def _on_render_thread(self) -> bool:
return self._render_thread is not None and threading.get_ident() == self._render_thread

Expand Down
Loading
Loading