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
19 changes: 19 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -576,6 +576,25 @@ policies are unchanged.
even when the mode name stays the same. The display republished its
current state only on a mode change or every 30 s, so `is_display_active`
and `on_demand_active` could be up to 30 s out of date.
- `/api/v3/display/current-status` no longer answers `mode: null` over the
control socket (#735) once the same mode has been on screen for more than
two minutes: a live game under live priority, Vegas, or a single plugin.
On ledpi it returned nulls in every sample for 90 minutes while the
display was live. The state stream's version leaves out the timestamps
that move on every publish, and a subscriber's keepalive tick carried only
the loop's heartbeat. So the web interface's copy kept the
`display.last_updated` of the last real change, and the reader's 120 s
rule called it unknown. The plugin runtime section had the same problem:
with no plugin changing state, `/plugins/state` and the `runtime` in
`/plugins/installed` read `stale` after 180 s. A tick (and a `state.get`
answer with `since`) now carries `volatile`: the current values of those
timestamps (`display.last_updated`, `on_demand.last_updated` and
`remaining`, `plugins.published_at`), and the subscription merges them
into its copy. The verdicts are unchanged. A render thread that stops
publishing still reads as `stalled` after 60 s and as unknown after 120 s,
a runtime publisher that stops still goes `stale`, and a subscription that
goes quiet still falls back to the cache. The cache path's 120 s rule is
unchanged.

### Scrolling

Expand Down
18 changes: 13 additions & 5 deletions docs/IPC_CONTROL_SOCKET.md
Original file line number Diff line number Diff line change
Expand Up @@ -146,7 +146,7 @@ connection until either side hangs up:

```json
{"v": 1, "id": "<the subscribe id>", "event": "state", "result": {...a state snapshot...}}
{"v": 1, "id": "<the subscribe id>", "event": "tick", "result": {"version": 7, "epoch": "…", "pid": 812, "served_at": 1790000000.1, "changed": false, "loop": {...}}}
{"v": 1, "id": "<the subscribe id>", "event": "tick", "result": {"version": 7, "epoch": "…", "pid": 812, "served_at": 1790000000.1, "changed": false, "loop": {...}, "volatile": {"display": {"last_updated": 1790000000.0}, "...": "..."}}}
```

An event has `event` where a response has `ok`, which is how a reader tells
Expand Down Expand Up @@ -216,8 +216,14 @@ socket.
counts within an `epoch`, one run of the display process, so a reader that
sees a new `epoch` has a restarted display.
- `state.get` with `since` and `epoch` from an earlier answer gets just
`{changed: false, version, epoch, pid, served_at, loop}` while nothing has
changed.
`{changed: false, version, epoch, pid, served_at, loop, volatile}` while
nothing has changed. `volatile` is `{section: {key: value}}`: the current
values of those ignored timestamps, which the reader merges into the copy
it has. They don't make a new version, but they are still news:
`display.last_updated` is how a reader knows the render thread is still
publishing, and `plugins.published_at` the runtime publisher. Without
them a reader's copy kept the timestamps of the last real change, so a
mode on screen for over 120 s read as unknown.
- A snapshot that would not fit in a message (hundreds of plugins) is sent
without `plugins`, and `truncated: ["plugins"]` says so. Readers then use
the cache for that section only.
Expand All @@ -226,8 +232,10 @@ socket.

- a `state` event (a full snapshot) whenever the version changes, and
- a `tick` at least every 5 s (`SUBSCRIBE_KEEPALIVE_SECONDS`) when nothing
changed. It carries `loop`, so a stalled render loop shows up within one
tick, and it tells the reader the connection is alive.
changed. It is the short `changed: false` answer, so it carries `loop`
(a stalled render loop shows up within one tick) and `volatile` (the
timestamps stay as fresh as the writers keep them), and it tells the
reader the connection is alive.

A slow reader is never sent a backlog: each event is the latest version, so
one that falls behind skips the versions in between. A reader that has heard
Expand Down
32 changes: 30 additions & 2 deletions src/ipc/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -283,6 +283,29 @@ def snapshot_loop_age(snapshot: Mapping[str, Any],
return max(float(age), 0.0) + snapshot_age(snapshot, now_mono)


def _merge_volatile(state: Dict[str, Any], volatile: Any) -> None:
"""Fold a tick's ``volatile`` values (``{section: {key: value}}``) into
``state``, copying each section it touches.

These are the timestamps the hub leaves out of its version --
``display.last_updated``, ``on_demand.last_updated``/``remaining``,
``plugins.published_at`` -- and the readers judge freshness by them, so
a copy that only full ``state`` events updated would go stale while the
same mode stayed on screen. Only keys the section already has are taken:
a tick never adds a section or a key the last snapshot did not carry
(a section left out of a truncated snapshot stays out).
"""
if not isinstance(volatile, dict):
return # a display from before ticks carried them
for name, values in volatile.items():
section = state.get(name)
if not isinstance(section, dict) or not isinstance(values, dict):
continue
fresh = {k: v for k, v in values.items() if k in section}
if fresh:
state[name] = dict(section, **fresh)


#: A subscription that has heard nothing for this long is not trusted: the
#: display sends a tick at least every SUBSCRIBE_KEEPALIVE_SECONDS.
SUBSCRIPTION_SILENCE_SECONDS = 3 * SUBSCRIBE_KEEPALIVE_SECONDS
Expand Down Expand Up @@ -449,12 +472,17 @@ def _store(self, result: Dict[str, Any], full: bool) -> None:
self.snapshots += 1
elif (self._snapshot is not None
and result.get('epoch') == self._snapshot.get('epoch')):
# A tick: nothing changed but the render loop's liveness.
# A tick: nothing changed but the render loop's liveness and
# the volatile keys (timestamps) the writers keep refreshing.
snap = dict(self._snapshot)
state = dict(snap.get('state') or {})
if result.get('version') == snap.get('version'):
_merge_volatile(state, result.get('volatile'))
loop = result.get('loop')
if isinstance(loop, dict):
snap['state'] = dict(snap.get('state') or {}, loop=loop)
state['loop'] = loop
snap['loop'] = loop
snap['state'] = state
snap['served_at'] = result.get('served_at', snap.get('served_at'))
self._snapshot = snap
else:
Expand Down
21 changes: 14 additions & 7 deletions src/ipc/contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,7 @@ class Command:

#: A subscriber hears from the display at least this often: a ``state``
#: event when something changed, else a ``tick`` carrying the render loop's
#: liveness. A client that has heard nothing for a few of these treats its
#: liveness and the latest volatile timestamps. A client that has heard nothing for a few of these treats its
#: copy as unknown.
SUBSCRIBE_KEEPALIVE_SECONDS = 5.0

Expand Down Expand Up @@ -530,8 +530,10 @@ class StateGetArgs:
"""``state.get``: the display's state, as a versioned snapshot.

With ``since`` and the ``epoch`` it came from, the answer is only
``{changed: false, version, epoch, served_at, loop}`` while the state is
still at that version, so a poller that already has it is sent no state.
``{changed: false, version, epoch, served_at, loop, volatile}`` while the
state is still at that version, so a poller that already has it is sent
no state -- only the latest values of the keys that do not count as a
change (``volatile``, see :class:`StateSnapshot`).
"""
since: Optional[int] = None
epoch: Optional[str] = None
Expand Down Expand Up @@ -672,9 +674,13 @@ class StateSnapshot(TypedDict, total=False):
``version`` counts changes to the state within one ``epoch`` (one run of
the display process): a reader that sees a new epoch starts over.
``changed`` is False only for a ``state.get`` whose ``since`` is still
current, and then ``state`` is absent. ``served_at`` is the display's
wall clock when it answered. ``loop`` is measured at that moment, so it
is also inside ``state``.
current, and then ``state`` is absent and ``volatile`` is there instead:
``{section: {key: value}}``, the current values of the keys the version
ignores (``display.last_updated``, ``on_demand.last_updated`` and
``remaining``, ``plugins.published_at``). A reader merges them into the
copy it has; they are how it can tell the writers are still publishing.
``served_at`` is the display's wall clock when it answered. ``loop`` is
measured at that moment, so it is also inside ``state``.

``state`` holds the sections in :data:`STATE_SECTIONS`:

Expand All @@ -696,12 +702,13 @@ class StateSnapshot(TypedDict, total=False):
served_at: float
changed: bool
state: Dict[str, Any]
volatile: Dict[str, Dict[str, Any]]
loop: LoopState


class StateEventKind:
STATE = 'state' # result: a full StateSnapshot, the latest version
TICK = 'tick' # result: {version, epoch, pid, served_at, loop}; nothing changed
TICK = 'tick' # result: {version, epoch, pid, served_at, loop, volatile}; nothing changed


@dataclass(frozen=True)
Expand Down
33 changes: 30 additions & 3 deletions src/ipc/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,22 @@ def _fingerprint(value: Optional[Mapping[str, Any]], volatile: Iterable[str]) ->
return {k: v for k, v in value.items() if k not in skip} if skip else dict(value)


def _volatile_values(sections: Mapping[str, Optional[Dict[str, Any]]],
volatile: Mapping[str, FrozenSet[str]]) -> Dict[str, Dict[str, Any]]:
"""``{section: {key: value}}``: the volatile keys each published section
has now. The sections are never mutated after publish (a publish swaps
in a new dict), so reading them outside the lock is safe."""
values: Dict[str, Dict[str, Any]] = {}
for name, keys in volatile.items():
value = sections.get(name)
if not keys or not isinstance(value, dict):
continue
present = {k: value[k] for k in keys if k in value}
if present:
values[name] = present
return values


def _unknown_loop() -> Dict[str, Any]:
return {'heartbeat_age_seconds': None, 'armed': False, 'stale_after': None}

Expand Down Expand Up @@ -259,7 +275,11 @@ class StateHub:
(``loop_probe``), so it keeps ageing while the render thread is stuck.

The version goes up when a section's value changes, ignoring the keys
the publisher names as volatile (timestamps). Publishing never blocks on
the publisher names as volatile (timestamps). Those keys still carry
news -- ``display.last_updated`` is the render thread's proof of life --
so the short ``changed: false`` answer, which is what a subscriber's
tick carries, has their current values in ``volatile``; a reader merges
them into its copy. Publishing never blocks on
a reader: the lock is held only to swap a dict reference and compare it,
and every socket write happens on the reader's own thread, outside it.
A reader that is slow gets the latest version when it next asks, not
Expand All @@ -274,6 +294,7 @@ def __init__(self, loop_probe: Optional[LoopProbe] = None, *,
self._cond = threading.Condition(threading.Lock())
self._sections: Dict[str, Optional[Dict[str, Any]]] = {}
self._fingerprints: Dict[str, Any] = {}
self._volatile: Dict[str, FrozenSet[str]] = {}
self._version = 0
self.epoch = epoch or uuid.uuid4().hex[:16]
self.pid = os.getpid() if pid is None else pid
Expand Down Expand Up @@ -301,9 +322,11 @@ def publish(self, section: str, value: Optional[Mapping[str, Any]],
The value is copied (one level), so the caller may reuse its dict.
"""
stored = None if value is None else dict(value)
fingerprint = _fingerprint(stored, volatile)
skip = frozenset(volatile)
fingerprint = _fingerprint(stored, skip)
with self._cond:
self._sections[section] = stored
self._volatile[section] = skip
if self._fingerprints.get(section, _MISSING) == fingerprint:
return False
self._fingerprints[section] = fingerprint
Expand Down Expand Up @@ -333,11 +356,14 @@ def snapshot(self, since: Optional[int] = None,
"""The :class:`~src.ipc.contract.StateSnapshot` now.

``since`` with this hub's ``epoch``, still the current version, gives
the short ``changed: false`` form.
the short ``changed: false`` form, with ``volatile``: each section's
volatile keys at their latest values (the rest of the section is
what the reader already has).
"""
with self._cond:
version = self._version
sections = dict(self._sections)
volatile = dict(self._volatile)
loop = self.loop()
result: Dict[str, Any] = {
'schema': STATE_SCHEMA,
Expand All @@ -349,6 +375,7 @@ def snapshot(self, since: Optional[int] = None,
}
if since is not None and epoch == self.epoch and since == version:
result['changed'] = False
result['volatile'] = _volatile_values(sections, volatile)
return result
state: Dict[str, Any] = {name: sections.get(name) for name in STATE_SECTIONS
if name != 'loop'}
Expand Down
82 changes: 82 additions & 0 deletions test/test_ipc_state_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,27 @@ def test_since_the_current_version_is_the_short_answer(self, hub):
assert hub.snapshot(since=1, epoch='other')['changed'] is True
assert hub.snapshot(since=1)['changed'] is True

def test_the_short_answer_carries_the_latest_volatile_values(self, hub):
"""A version that stays put must not freeze the timestamps a reader
judges freshness by: the short answer (a tick) carries them."""
hub.publish('on_demand', {'active': False, 'last_updated': 2.0, 'remaining': 0.0},
volatile=('last_updated', 'remaining'))
hub.publish('brightness', {'brightness': 50})
version = hub.version
assert not hub.publish('display', _display(updated=1234.5), volatile=('last_updated',))
short = hub.snapshot(since=version, epoch='e1')
assert short['changed'] is False and 'state' not in short
assert short['volatile'] == {'display': {'last_updated': 1234.5},
'on_demand': {'last_updated': 2.0, 'remaining': 0.0}}
# The full answer has them in place already.
assert 'volatile' not in hub.snapshot()

def test_volatile_values_skip_sections_without_any(self):
hub = StateHub(epoch='e1')
hub.publish('brightness', {'brightness': 50})
hub.publish('plugins', None, volatile=('published_at',))
assert hub.snapshot(since=hub.version, epoch='e1')['volatile'] == {}

def test_each_display_run_has_its_own_epoch(self):
assert StateHub().epoch != StateHub().epoch

Expand Down Expand Up @@ -276,6 +297,40 @@ def test_a_tick_refreshes_the_loop_and_keeps_the_state(self, hub):
assert latest['state']['loop']['heartbeat_age_seconds'] == 70.0
assert latest['received_mono'] == clock.now

def test_a_tick_refreshes_the_volatile_timestamps(self, hub):
"""The bug from the ledpi rig: the same mode on screen for minutes
left the reader's display.last_updated at the last real change."""
clock = FakeClock()
sub = client.StateSubscription(paths=['/nowhere'], clock=clock)
sub._store(hub.snapshot(), full=True)
hub.publish('display', _display(updated=9999.0), volatile=('last_updated',))
sub._store(hub.snapshot(since=1, epoch='e1'), full=False)
latest = sub.latest()
assert latest['state']['display']['last_updated'] == 9999.0
assert latest['state']['display']['mode'] == 'clock'
assert latest['version'] == 1

def test_a_tick_adds_no_section_or_key_and_ignores_another_version(self, hub):
clock = FakeClock()
sub = client.StateSubscription(paths=['/nowhere'], clock=clock)
snap = hub.snapshot()
sub._store(snap, full=True)
updated = snap['state']['display']['last_updated']
sub._store({'version': 1, 'epoch': 'e1',
'volatile': {'display': {'new_key': 1}, 'plugins': {'published_at': 5.0},
'on_demand': 'junk'}}, full=False)
state = sub.latest()['state']
assert 'new_key' not in state['display'] and state['plugins'] is None
assert state['on_demand'] is None
# A tick for a version this copy is not at says nothing about it.
sub._store({'version': 7, 'epoch': 'e1',
'volatile': {'display': {'last_updated': 5.0}}}, full=False)
assert sub.latest()['state']['display']['last_updated'] == updated
# A display from before ticks carried them: the loop still refreshes.
sub._store({'version': 1, 'epoch': 'e1', 'loop': {'heartbeat_age_seconds': 3.0}},
full=False)
assert sub.latest()['loop'] == {'heartbeat_age_seconds': 3.0}

def test_a_tick_from_another_epoch_is_ignored(self, hub):
clock = FakeClock()
sub = client.StateSubscription(paths=['/nowhere'], clock=clock)
Expand Down Expand Up @@ -369,6 +424,33 @@ def test_ticks_keep_a_quiet_subscription_alive(self, live, path):
assert kinds == ['tick', 'tick', 'tick']
sock.close()

def test_ticks_carry_the_timestamps_the_version_ignores(self, live, path):
_, hub = live
sock, reader = self._subscribe(path)
pending = []
self._next(sock, reader, pending)
hub.publish('display', _display(updated=4242.0), volatile=('last_updated',))
for _ in range(3): # a tick already on its way may predate the publish
tick = self._next(sock, reader, pending)
assert tick['event'] == 'tick'
if tick['result']['volatile'] == {'display': {'last_updated': 4242.0}}:
break
else:
pytest.fail(f'no tick carried the new last_updated: {tick}')
sock.close()

def test_a_subscription_keeps_last_updated_current(self, live, path):
_, hub = live
sub = client.StateSubscription(paths=[path]).start()
try:
assert _wait_until(lambda: sub.latest() is not None)
hub.publish('display', _display(updated=4242.0), volatile=('last_updated',))
assert _wait_until(lambda: (sub.latest() or {}).get('state', {})
.get('display', {}).get('last_updated') == 4242.0)
assert sub.snapshots == 1 # ticks, not new snapshots
finally:
sub.stop()

def test_a_burst_is_coalesced_to_the_latest(self, live, path):
_, hub = live
sock, reader = self._subscribe(path)
Expand Down
Loading
Loading