-
-
Notifications
You must be signed in to change notification settings - Fork 29
Expand file tree
/
Copy pathclient.py
More file actions
557 lines (475 loc) · 23.9 KB
/
Copy pathclient.py
File metadata and controls
557 lines (475 loc) · 23.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
"""The web side of the control socket: one request, a short timeout, no retries.
Every failure -- no socket (the display is stopped, or predates the socket),
a refused or timed-out connection, a reply that breaks the contract, or an
error the display returned -- raises :class:`ControlError` with a short
``reason``. Nothing here blocks for longer than ``timeout`` in total.
Whether the caller may then write the file mailbox instead is
:func:`should_fall_back`: only when the display never took the request (it
could not be reached, or it is too old to know the command). A display that
took the request and then failed, refused or went quiet is answered as
that, not posted a second time through the mailbox.
"""
from __future__ import annotations
import socket
import threading
import time
import uuid
from typing import Any, Callable, Dict, List, Mapping, Optional, Sequence
from src.ipc.contract import (
AWAIT_SECONDS,
MAX_MESSAGE_BYTES,
PROTOCOL_VERSION,
SUBSCRIBE_KEEPALIVE_SECONDS,
SUPPORTED_VERSIONS,
Command,
ErrorCode,
FrameReader,
ProtocolError,
Request,
Response,
StateEvent,
StateEventKind,
client_socket_paths,
decode_message,
encode_message,
parse_args,
socket_supported,
)
#: Total budget for one request: connect, send and the reply. The display
#: answers from a thread that does no rendering, normally within a few
#: milliseconds; this only bounds a wedged one. The web route then falls back
#: to the mailbox, so a timeout costs this much latency and nothing else.
DEFAULT_TIMEOUT_SECONDS = 1.0
class ControlError(Exception):
"""The socket could not carry the request. ``reason`` is a short code.
Transport reasons: ``disabled``, ``unsupported``, ``no_socket``,
``refused``, ``timeout``, ``closed``, ``bad_response``, ``invalid_request``.
When the display answered with an error, ``reason`` is that error's
:class:`~src.ipc.contract.ErrorCode` (``busy``, ``unknown_command``, ...).
``sent`` is True once the whole request was written to a connected
display, which may then have acted on it. A refusal the display sends
before it reads anything (``forbidden``, too many connections) carries
no request id and leaves ``sent`` False.
"""
def __init__(self, reason: str, message: str = '', *, sent: bool = False):
super().__init__(reason, message)
self.reason = reason
self.message = message
self.sent = sent
def __str__(self) -> str:
return f'{self.reason}: {self.message}' if self.message else self.reason
#: Answers from a display that read the request but does not speak it: one
#: older than the command (an upgrade in progress) or the protocol version.
#: It did nothing, so the mailbox is the way to reach it.
UPGRADE_REASONS = frozenset({ErrorCode.UNKNOWN_COMMAND, ErrorCode.UNSUPPORTED_VERSION})
def should_fall_back(error: BaseException) -> bool:
"""May the caller write the file mailbox after ``error``?
Yes when the display never took the request: there is no socket (the
display is stopped, predates the socket, or it is switched off), the
connection was refused or timed out, the display turned the connection
away before reading it, or it is too old to know the command
(:data:`UPGRADE_REASONS`). Also for an error that is not a
:class:`ControlError` (a bug in the client), as before.
No once the display had the request: a ``busy`` queue, ``invalid_args``,
an ``internal`` error, or a timeout or hang-up after the request was
sent. The display may have applied it, or would refuse it from the
mailbox too, so a second copy there only hides the failure.
"""
if not isinstance(error, ControlError):
return True
return not error.sent or error.reason in UPGRADE_REASONS
def request(cmd: str, args: Optional[Mapping[str, Any]] = None, *,
request_id: Optional[str] = None,
timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""Send one command and return its ``result``. Raises :class:`ControlError`."""
args = dict(args or {})
request_id = request_id or str(uuid.uuid4())
try:
# Refuse locally what the display would refuse: a malformed id
# (callers may pass their own) or arguments that break the contract.
envelope = Request.from_dict({'v': PROTOCOL_VERSION, 'id': request_id,
'cmd': cmd, 'args': args})
parse_args(cmd, args)
payload = encode_message(envelope.to_dict())
except ProtocolError as e:
raise ControlError('invalid_request', e.message) from None
if not socket_supported():
raise ControlError('unsupported', 'no Unix sockets on this platform')
candidates: List[str] = list(paths) if paths is not None else client_socket_paths()
if not candidates:
raise ControlError('disabled', 'the control socket is turned off')
deadline = time.monotonic() + timeout
sock = _connect(candidates, deadline)
try:
response = _exchange(sock, payload, deadline)
finally:
sock.close()
# A refusal before the request was read (forbidden, too many
# connections) carries no id.
if response.id != request_id and not (response.id is None and not response.ok):
raise ControlError('bad_response', 'the reply is for a different request', sent=True)
if not response.ok:
error = response.error
# No id: refused at the door (forbidden, too many connections),
# before the display read the request.
raise ControlError(error.code if error else 'bad_response',
error.message if error else '',
sent=response.id is not None)
return dict(response.result or {})
def _remaining(deadline: float) -> float:
left = deadline - time.monotonic()
if left <= 0:
raise ControlError('timeout', 'no reply in time')
return left
def _connect(paths: Sequence[str], deadline: float) -> socket.socket:
last = ControlError('no_socket', 'the display is not serving the control socket')
for path in paths:
sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
try:
sock.settimeout(_remaining(deadline))
sock.connect(path)
return sock
except (FileNotFoundError, NotADirectoryError):
sock.close()
continue
except ConnectionRefusedError:
sock.close()
last = ControlError('refused', f'nothing is listening at {path}')
except BlockingIOError:
# EAGAIN: the listen backlog is full -- a live but swamped display.
sock.close()
raise ControlError('busy', 'the display is not accepting connections') from None
except PermissionError:
sock.close()
last = ControlError('refused', f'no permission to connect to {path}')
except socket.timeout:
sock.close()
raise ControlError('timeout', 'connect timed out') from None
except ControlError:
sock.close()
raise
except OSError as e:
sock.close()
last = ControlError('refused', f'{path}: {e}')
raise last
def _exchange(sock: socket.socket, payload: bytes, deadline: float) -> Response:
"""Send ``payload`` and read the reply. A failure once the whole request
is written raises with ``sent=True``: the display may have it."""
sent = False
try:
sock.settimeout(_remaining(deadline))
sock.sendall(payload)
sent = True
reader = FrameReader(MAX_MESSAGE_BYTES)
while True:
sock.settimeout(_remaining(deadline))
data = sock.recv(4096)
if not data:
raise ControlError('closed', 'the display closed the connection', sent=sent)
lines = reader.feed(data)
if lines:
return Response.from_dict(decode_message(lines[0]))
except socket.timeout:
raise ControlError('timeout', 'no reply in time', sent=sent) from None
except ProtocolError as e:
raise ControlError('bad_response', e.message, sent=sent) from None
except ControlError as e:
e.sent = e.sent or sent
raise
except OSError as e:
raise ControlError('closed', str(e), sent=sent) from None
# -- commands ---------------------------------------------------------------------------
def on_demand_start(request_id: str, plugin_id: Optional[str], mode: Optional[str],
duration: Any = None, pinned: bool = False, *,
timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""Ask the display to show a plugin now. Returns the ack; raises :class:`ControlError`.
``request_id`` doubles as the on-demand request id, so a request that a
timed-out caller then also writes to the mailbox is processed only once.
"""
args = {'plugin_id': plugin_id, 'mode': mode, 'duration': duration, 'pinned': pinned}
return request(Command.ON_DEMAND_START, args, request_id=request_id,
timeout=timeout, paths=paths)
def on_demand_stop(request_id: str, *, timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""Ask the display to end on-demand. Returns the ack; raises :class:`ControlError`."""
return request(Command.ON_DEMAND_STOP, {}, request_id=request_id,
timeout=timeout, paths=paths)
def on_demand_status(*, timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""The display's live on-demand state. Raises :class:`ControlError`."""
return request(Command.ON_DEMAND_STATUS, {}, timeout=timeout, paths=paths)
#: Headroom over the display's own wait for an awaited command, so its
#: ``pending`` answer arrives before the client gives up.
_AWAIT_MARGIN_SECONDS = 1.0
def _awaited_timeout(cmd: str) -> float:
return AWAIT_SECONDS[cmd] + _AWAIT_MARGIN_SECONDS
def brightness_set(brightness: int, *, timeout: Optional[float] = None,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""Set the panel's normal brightness now (transient: config.json is not
written). Returns the applied :class:`~src.ipc.contract.BrightnessResult`;
raises :class:`ControlError`.
"""
return request(Command.BRIGHTNESS_SET, {'brightness': brightness},
timeout=_awaited_timeout(Command.BRIGHTNESS_SET) if timeout is None
else timeout, paths=paths)
def plugin_reload(plugin_id: str, *, timeout: Optional[float] = None,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""Have the display reload a running plugin from disk.
Returns :class:`~src.ipc.contract.PluginReloadResult` once the new code is
running. Raises :class:`ControlError`: ``not_loaded`` (not running it),
``failed`` (the new version did not load), ``pending`` (not done in
time; it will still happen), or a transport reason.
"""
return request(Command.PLUGIN_RELOAD, {'plugin_id': plugin_id},
timeout=_awaited_timeout(Command.PLUGIN_RELOAD) if timeout is None
else timeout, paths=paths)
def errors_clear(request_id: str, cutoff: float, *,
timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""Have the display forget the plugin errors recorded at or before
``cutoff`` (epoch seconds) and publish its error snapshot again.
Returns :class:`~src.ipc.contract.ErrorsClearResult` once it is done.
Raises :class:`ControlError`: ``unknown_command`` from a display older
than the command, which still reads the ``plugin_error_clear_request``
mailbox.
"""
return request(Command.ERRORS_CLEAR, {'cutoff': cutoff}, request_id=request_id,
timeout=timeout, paths=paths)
def ping(*, timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
return request(Command.PING, {}, timeout=timeout, paths=paths)
def hello(client: str = 'web', *, timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""Version negotiation: the result's ``version`` is the one both sides speak."""
return request(Command.HELLO, {'versions': list(SUPPORTED_VERSIONS), 'client': client},
timeout=timeout, paths=paths)
# -- the state stream (stage 3) ---------------------------------------------------------
def state_get(since: Optional[int] = None, epoch: Optional[str] = None, *,
timeout: float = DEFAULT_TIMEOUT_SECONDS,
paths: Optional[Sequence[str]] = None) -> Dict[str, Any]:
"""The display's state now, as a :class:`~src.ipc.contract.StateSnapshot`.
With ``since``/``epoch`` from an earlier answer, an unchanged state comes
back in the short ``changed: false`` form. Raises :class:`ControlError`
(``unknown_command`` from a display older than stage 3).
"""
args: Dict[str, Any] = {}
if since is not None:
args['since'] = since
if epoch is not None:
args['epoch'] = epoch
return request(Command.STATE_GET, args, timeout=timeout, paths=paths)
def snapshot_age(snapshot: Mapping[str, Any], now_mono: Optional[float] = None) -> float:
"""Seconds since ``snapshot`` arrived: ``received_mono`` (set by
:meth:`StateSubscription.latest`) to now; 0 for a one-shot answer."""
received = snapshot.get('received_mono')
if isinstance(received, (int, float)) and not isinstance(received, bool):
now_mono = time.monotonic() if now_mono is None else now_mono
return max(now_mono - float(received), 0.0)
return 0.0
def snapshot_loop_age(snapshot: Mapping[str, Any],
now_mono: Optional[float] = None) -> Optional[float]:
"""The render loop's heartbeat age now, from a state snapshot: the age the
display measured when it answered, plus the time since the answer
arrived. None when the display has no beat to report yet."""
loop = snapshot.get('loop')
if not isinstance(loop, dict):
state = snapshot.get('state')
loop = state.get('loop') if isinstance(state, dict) else None
age = loop.get('heartbeat_age_seconds') if isinstance(loop, dict) else None
if not isinstance(age, (int, float)) or isinstance(age, bool):
return None
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
#: Reconnect backoff: the first retry, and the cap. A display that does not
#: know state.subscribe (stage 2 or older) is retried at the cap.
_RECONNECT_MIN_SECONDS = 1.0
_RECONNECT_MAX_SECONDS = 30.0
#: Failures that another try soon will not fix.
_SLOW_RETRY_REASONS = frozenset({'unknown_command', 'unsupported_version', 'disabled',
'unsupported'})
class StateSubscription:
"""One ``state.subscribe`` connection, held on a daemon thread.
Keeps the latest snapshot the display pushed, so a reader answers from
memory (:meth:`latest`). Reconnects with a backoff when the display goes
away. Never raises into the caller: :meth:`latest` is None whenever the
copy cannot be vouched for (not connected, or silent for longer than
``silence``), and the caller falls back.
"""
def __init__(self, paths: Optional[Sequence[str]] = None, *,
silence: float = SUBSCRIPTION_SILENCE_SECONDS,
connect_timeout: float = DEFAULT_TIMEOUT_SECONDS,
clock: Callable[[], float] = time.monotonic):
self._paths = list(paths) if paths is not None else None
self._silence = silence
self._connect_timeout = connect_timeout
self._clock = clock
self._lock = threading.Lock()
self._snapshot: Optional[Dict[str, Any]] = None
self._received: Optional[float] = None
self._connected = False
self._stop = threading.Event()
self._sock: Optional[socket.socket] = None
self._thread: Optional[threading.Thread] = None
#: The reason the last connection ended (a ControlError reason).
self.last_error: Optional[str] = None
#: Full snapshots received: the subscribe answer and each state event.
self.snapshots = 0
# -- the reader's side ---------------------------------------------------
@property
def connected(self) -> bool:
return self._connected
def latest(self) -> Optional[Dict[str, Any]]:
"""A copy of the latest snapshot, with ``received_mono`` (this
process's monotonic clock when it arrived); None when not trusted."""
with self._lock:
if not self._connected or self._snapshot is None or self._received is None:
return None
if self._clock() - self._received > self._silence:
return None
snap = dict(self._snapshot)
snap['received_mono'] = self._received
return snap
# -- lifecycle -----------------------------------------------------------
def start(self) -> 'StateSubscription':
if self._thread is None or not self._thread.is_alive():
self._stop.clear()
self._thread = threading.Thread(target=self._run, name='ledmatrix-state-feed',
daemon=True)
self._thread.start()
return self
def stop(self, timeout: float = 2.0) -> None:
self._stop.set()
sock = self._sock
if sock is not None:
try:
sock.shutdown(socket.SHUT_RDWR)
except OSError:
pass
thread = self._thread
if thread is not None and thread is not threading.current_thread():
thread.join(timeout)
self._thread = None
# -- the feed thread -----------------------------------------------------
def _run(self) -> None:
backoff = _RECONNECT_MIN_SECONDS
while not self._stop.is_set():
snapshots = self.snapshots
try:
self._follow()
except ControlError as e:
self.last_error = e.reason
if e.reason in _SLOW_RETRY_REASONS:
backoff = _RECONNECT_MAX_SECONDS
except Exception as e: # pylint: disable=broad-except
self.last_error = type(e).__name__
finally:
with self._lock:
self._connected = False
sock, self._sock = self._sock, None
if sock is not None:
try:
sock.close()
except OSError:
pass
if self.snapshots != snapshots:
# This connection got as far as the display's state: whatever
# ended it (a restart, most often), it was working, so the
# next try starts from the shortest wait again.
backoff = _RECONNECT_MIN_SECONDS
if self._stop.wait(backoff):
return
backoff = min(backoff * 2, _RECONNECT_MAX_SECONDS)
def _follow(self) -> None:
"""Subscribe, then read events until the connection ends. Raises ControlError."""
if not socket_supported():
raise ControlError('unsupported', 'no Unix sockets on this platform')
candidates = list(self._paths) if self._paths is not None else client_socket_paths()
if not candidates:
raise ControlError('disabled', 'the control socket is turned off')
request_id = str(uuid.uuid4())
payload = encode_message(Request(id=request_id, cmd=Command.STATE_SUBSCRIBE,
args={}).to_dict())
sock = _connect(candidates, time.monotonic() + self._connect_timeout)
self._sock = sock
try:
sock.settimeout(self._connect_timeout)
sock.sendall(payload)
# A read waits for the next event; the display sends one at least
# every keepalive, so this much silence means it is gone.
sock.settimeout(self._silence)
reader = FrameReader(MAX_MESSAGE_BYTES)
first = True
while not self._stop.is_set():
data = sock.recv(65536)
if not data:
raise ControlError('closed', 'the display closed the connection')
for line in reader.feed(data):
obj = decode_message(line)
if first:
response = Response.from_dict(obj)
if not response.ok:
error = response.error
raise ControlError(error.code if error else 'bad_response',
error.message if error else '')
self._store(dict(response.result or {}), full=True)
first = False
continue
event = StateEvent.from_dict(obj)
self._store(event.result, full=event.event == StateEventKind.STATE)
except socket.timeout:
raise ControlError('timeout', 'the display went quiet') from None
except ProtocolError as e:
raise ControlError('bad_response', e.message) from None
except OSError as e:
if self._stop.is_set():
return
raise ControlError('closed', str(e)) from None
def _store(self, result: Dict[str, Any], full: bool) -> None:
now = self._clock()
with self._lock:
if full and isinstance(result.get('state'), dict):
self._snapshot = result
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 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):
state['loop'] = loop
snap['loop'] = loop
snap['state'] = state
snap['served_at'] = result.get('served_at', snap.get('served_at'))
self._snapshot = snap
else:
return # a tick before any state, or from another epoch
self._received = now
self._connected = True