Coverage for src/signalk_cli/stream/api.py: 100%

125 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-10-06 00:17 +0000

1"""Python API for the SignalK v1 Streaming (delta) API. 

2 

3Spec: https://signalk.org/specification/1.8.2/doc/streaming_api.html 

4Subscribe path wildcards: https://signalk.org/specification/1.8.2/doc/subscription_protocol.html 

5 

6Examples: 

7 ```python 

8 from signalk_cli.stream import StreamClient 

9 client = StreamClient("http://boat.local:3000") 

10 table = client.collect(["navigation.*"], count=100, policy="instant") 

11 import polars as pl 

12 df = pl.DataFrame(table) 

13 ``` 

14""" 

15 

16import fnmatch 

17import json 

18from collections.abc import Iterator, Sequence 

19from dataclasses import dataclass 

20from typing import Any, Literal, NamedTuple, Self 

21from urllib.parse import urlparse, urlunparse 

22 

23import niquests 

24 

25from .._arrow import ArrowTable 

26from ..errors import SignalKError 

27from ..net import normalise_host 

28 

29STREAM_PATH = "/signalk/v1/stream" 

30 

31SUBSCRIBE_POLICIES = ["none", "self", "all"] 

32SUBSCRIPTION_POLICIES = ["instant", "ideal", "fixed"] 

33 

34Subscribe = Literal["none", "self", "all"] 

35Policy = Literal["instant", "ideal", "fixed"] 

36 

37 

38def to_ws_url(host: str) -> str: 

39 """Convert an http(s) host base URL to the ws(s) streaming endpoint URL.""" 

40 parsed = urlparse(host) 

41 scheme = "wss" if parsed.scheme == "https" else "ws" 

42 return urlunparse((scheme, parsed.netloc, STREAM_PATH, "", "", "")) 

43 

44 

45def build_subscribe_message( 

46 context: str, 

47 paths: Sequence[str], 

48 *, 

49 period_ms: int | None = None, 

50 policy: str | None = None, 

51 min_period_ms: int | None = None, 

52) -> dict: 

53 """Build a client subscribe message for the given context and paths. 

54 

55 An empty path list subscribes to all paths within the context (equivalent 

56 to a single "*" path). Paths are passed through unchanged — wildcarding is 

57 handled server-side per the SignalK Subscription Protocol: "*" at the end 

58 of a path matches any suffix (e.g. "navigation.*"), and "*" as a middle 

59 segment matches any single segment there (e.g. "propulsion.*.oilTemperature"). 

60 

61 period_ms/policy/min_period_ms map directly to the per-path "period", 

62 "policy", and "minPeriod" fields of the Subscription Protocol and are 

63 applied identically to every path when given. `policy` 

64 ("instant"/"ideal"/"fixed") defaults to "ideal" server-side if omitted; 

65 `min_period_ms` only affects the "instant" policy. The protocol also 

66 defines a per-path "format" ("delta"/"full") field, but it's omitted 

67 here: signalk-server rejects "full" outright and always sends delta 

68 messages regardless, so exposing the choice would be misleading. 

69 """ 

70 entry_extra: dict[str, Any] = {} 

71 if period_ms is not None: 

72 entry_extra["period"] = period_ms 

73 if policy is not None: 

74 entry_extra["policy"] = policy 

75 if min_period_ms is not None: 

76 entry_extra["minPeriod"] = min_period_ms 

77 

78 path_list = list(paths) if paths else ["*"] 

79 return { 

80 "context": context, 

81 "subscribe": [{"path": p, **entry_extra} for p in path_list], 

82 } 

83 

84 

85# --------------------------------------------------------------------------- 

86# Messages and rows 

87# --------------------------------------------------------------------------- 

88 

89 

90def source_matches(source: str, patterns: Sequence[str]) -> bool: 

91 """Match a `$source` string against source filter patterns (OR'd). 

92 

93 No patterns means no filtering (always matches). A pattern containing 

94 glob metacharacters (`*`/`?`/`[`) is matched as-is via `fnmatch`; 

95 otherwise it's treated as a substring match, e.g. "Teltonika" matches 

96 the source "Teltonika.GP". 

97 """ 

98 if not patterns: 

99 return True 

100 return any( 

101 fnmatch.fnmatch(source, p if any(c in p for c in "*?[") else f"*{p}*") 

102 for p in patterns 

103 ) 

104 

105 

106def _update_source(update: dict) -> str: 

107 return update.get("$source") or json.dumps(update.get("source", {})) 

108 

109 

110class DeltaRow(NamedTuple): 

111 """One value (or meta entry) from a delta message. 

112 

113 `value` keeps its JSON type: a number, string, bool, object, array, or None. 

114 `kind` is `"value"`, or `"meta"` for metadata entries (units, 

115 description, zones, ...). 

116 """ 

117 

118 timestamp: str 

119 context: str 

120 source: str 

121 path: str 

122 kind: str 

123 value: Any 

124 

125 

126@dataclass(frozen=True) 

127class DeltaMessage: 

128 """A delta message from the server. 

129 

130 Attributes: 

131 text: The message exactly as received. 

132 payload: The decoded message. 

133 """ 

134 

135 text: str 

136 payload: dict 

137 

138 def rows( 

139 self, *, include_meta: bool = False, sources: Sequence[str] = () 

140 ) -> list[DeltaRow]: 

141 """Flatten into rows, one per value. 

142 

143 Args: 

144 include_meta: Also include each update's "meta" entries, as rows 

145 with `kind="meta"`. 

146 sources: Only include updates whose `$source` matches one of 

147 these patterns (see [`source_matches()`][signalk_cli.stream.api.source_matches]). 

148 """ 

149 context = self.payload.get("context", "") 

150 rows: list[DeltaRow] = [] 

151 for update in self.payload.get("updates", []): 

152 source = _update_source(update) 

153 if not source_matches(source, sources): 

154 continue 

155 timestamp = update.get("timestamp", "") 

156 kinds = ["values", "meta"] if include_meta else ["values"] 

157 for key in kinds: 

158 for entry in update.get(key, []): 

159 rows.append( 

160 DeltaRow( 

161 timestamp, 

162 context, 

163 source, 

164 entry.get("path", ""), 

165 "value" if key == "values" else "meta", 

166 entry.get("value"), 

167 ) 

168 ) 

169 return rows 

170 

171 def matches_sources(self, sources: Sequence[str]) -> bool: 

172 """True if any update in the message comes from a matching `$source`.""" 

173 return not sources or any( 

174 source_matches(_update_source(u), sources) 

175 for u in self.payload.get("updates", []) 

176 ) 

177 

178 

179def rows_to_arrow( 

180 rows: Sequence[DeltaRow], *, include_meta: bool = False 

181) -> ArrowTable: 

182 """Convert delta rows to a table: UTC `timestamp`, `context`, `source`, 

183 `path`, `kind` (only with `include_meta`) and `value`. 

184 

185 `value` is float64 if every value is a number, otherwise text, with 

186 objects and arrays as JSON. 

187 """ 

188 text = ["context", "source", "path", "kind"] 

189 columns: dict[str, list] = { 

190 "context": [r.context for r in rows], 

191 "source": [r.source for r in rows], 

192 "path": [r.path for r in rows], 

193 } 

194 if include_meta: 

195 columns["kind"] = [r.kind for r in rows] 

196 columns["value"] = [r.value for r in rows] 

197 return ArrowTable.from_rows( 

198 [r.timestamp or None for r in rows], columns, text_columns=text 

199 ) 

200 

201 

202# --------------------------------------------------------------------------- 

203# Stream and client 

204# --------------------------------------------------------------------------- 

205 

206 

207class DeltaStream: 

208 """An open subscription. Iterate it for [`DeltaMessage`][signalk_cli.stream.api.DeltaMessage] objects. 

209 

210 Use as a context manager, or call [`close()`][signalk_cli.stream.api.DeltaStream.close], to close the connection. 

211 

212 Raises: 

213 SignalKError: While iterating, if the connection is lost. 

214 """ 

215 

216 def __init__(self, ws: Any) -> None: 

217 self._ws = ws 

218 

219 def __enter__(self) -> Self: 

220 return self 

221 

222 def __exit__(self, *exc: object) -> None: 

223 self.close() 

224 

225 def close(self) -> None: 

226 """Close the WebSocket connection.""" 

227 self._ws.close() 

228 

229 def __iter__(self) -> Iterator[DeltaMessage]: 

230 return self.messages() 

231 

232 def messages(self, count: int | None = None) -> Iterator[DeltaMessage]: 

233 """Yield delta messages, skipping control messages such as the server's hello. 

234 

235 Stops after `count` messages if given, or when the server closes 

236 the connection. 

237 """ 

238 yielded = 0 

239 while count is None or yielded < count: 

240 try: 

241 payload = self._ws.next_payload() 

242 except niquests.RequestException as e: 

243 raise SignalKError.from_request(e) from e 

244 if payload is None: 

245 return 

246 if isinstance(payload, bytes): 

247 payload = payload.decode("utf-8", errors="replace") 

248 try: 

249 message = json.loads(payload) 

250 except (TypeError, json.JSONDecodeError): 

251 continue 

252 if not isinstance(message, dict) or "updates" not in message: 

253 continue 

254 yield DeltaMessage(payload, message) 

255 yielded += 1 

256 

257 def rows( 

258 self, 

259 count: int | None = None, 

260 *, 

261 include_meta: bool = False, 

262 sources: Sequence[str] = (), 

263 ) -> Iterator[DeltaRow]: 

264 """Yield rows from the next `count` messages (see [`rows()`][signalk_cli.stream.api.DeltaMessage.rows]).""" 

265 for message in self.messages(count): 

266 yield from message.rows(include_meta=include_meta, sources=sources) 

267 

268 def collect( 

269 self, 

270 count: int | None = None, 

271 *, 

272 include_meta: bool = False, 

273 sources: Sequence[str] = (), 

274 ) -> ArrowTable: 

275 """Read `count` messages (or until the server closes) into a table. 

276 

277 See [`rows_to_arrow()`][signalk_cli.stream.api.rows_to_arrow] for the columns. 

278 """ 

279 rows = list(self.rows(count, include_meta=include_meta, sources=sources)) 

280 return rows_to_arrow(rows, include_meta=include_meta) 

281 

282 

283class StreamClient: 

284 """Client for a SignalK server's v1 Streaming (delta) API. 

285 

286 Args: 

287 host: Server URL, e.g. `http://boat.local:3000` (`http://` is 

288 added if there's no scheme). 

289 context: SignalK context to subscribe to. Accepts the wildcard 

290 `*` (or `vessels.*`) for every vessel. 

291 subscribe: The connection-level auto-subscription the server adds at 

292 its own default rate: `"none"` (default), `"self"` or 

293 `"all"`. It's in addition to the explicit subscription for 

294 `context`; `"none"` avoids receiving your own vessel twice. 

295 session: A niquests session to connect with. One is created (and 

296 closed with the client) if not given. 

297 """ 

298 

299 def __init__( 

300 self, 

301 host: str, 

302 *, 

303 context: str = "vessels.self", 

304 subscribe: Subscribe = "none", 

305 session: niquests.Session | None = None, 

306 ) -> None: 

307 self.host = normalise_host(host).rstrip("/") 

308 self.context = context 

309 self.subscribe = subscribe 

310 self._owns_session = session is None 

311 self._session = session if session is not None else niquests.Session() 

312 

313 def __enter__(self) -> Self: 

314 return self 

315 

316 def __exit__(self, *exc: object) -> None: 

317 self.close() 

318 

319 def close(self) -> None: 

320 """Close the session, if the client created it.""" 

321 if self._owns_session: 

322 self._session.close() 

323 

324 def open( 

325 self, 

326 paths: Sequence[str] = (), 

327 *, 

328 policy: Policy | None = "ideal", 

329 period: float | None = 60.0, 

330 min_period: float | None = None, 

331 timeout: float | None = 30, 

332 ) -> DeltaStream: 

333 """Connect and subscribe, returning the open [`DeltaStream`][signalk_cli.stream.api.DeltaStream]. 

334 

335 Args: 

336 paths: Paths to subscribe to; all paths if empty. `*` wildcards 

337 are matched by the server, at the end of a path 

338 (`navigation.*`) or as a whole segment 

339 (`propulsion.*.oilTemperature`). 

340 policy: `"instant"` sends every change (limited by 

341 `min_period`); `"ideal"` also resends the last value if 

342 nothing changes within `period`; `"fixed"` sends the last 

343 value every `period`. 

344 period: Resend interval in seconds for `ideal`/`fixed`. 

345 min_period: Fastest send rate in seconds, for `instant`. 

346 timeout: Seconds to wait for each message, or None to wait 

347 indefinitely (a quiet subscription may send nothing for a 

348 long time). 

349 

350 Raises: 

351 SignalKError: If the connection fails. 

352 """ 

353 try: 

354 resp = self._session.get( 

355 to_ws_url(self.host), 

356 params={"subscribe": self.subscribe}, 

357 timeout=timeout, 

358 ) 

359 resp.raise_for_status() 

360 except niquests.RequestException as e: 

361 raise SignalKError.from_request(e) from e 

362 ws = resp.extension 

363 if ws is None: 

364 raise SignalKError( 

365 f"{self.host} did not accept a WebSocket connection " 

366 "(is niquests installed with the 'ws' extra?)" 

367 ) 

368 message = build_subscribe_message( 

369 self.context, 

370 paths, 

371 period_ms=None if period is None else int(period * 1000), 

372 policy=policy, 

373 min_period_ms=None if min_period is None else int(min_period * 1000), 

374 ) 

375 ws.send_payload(json.dumps(message)) 

376 return DeltaStream(ws) 

377 

378 def collect( 

379 self, 

380 paths: Sequence[str] = (), 

381 *, 

382 count: int, 

383 include_meta: bool = False, 

384 sources: Sequence[str] = (), 

385 policy: Policy | None = "ideal", 

386 period: float | None = 60.0, 

387 min_period: float | None = None, 

388 timeout: float | None = 30, 

389 ) -> ArrowTable: 

390 """Subscribe, read `count` delta messages into a table, and disconnect. 

391 

392 Arguments are as for [`open()`][signalk_cli.stream.api.StreamClient.open] and [`collect()`][signalk_cli.stream.api.DeltaStream.collect]. 

393 """ 

394 with self.open( 

395 paths, 

396 policy=policy, 

397 period=period, 

398 min_period=min_period, 

399 timeout=timeout, 

400 ) as stream: 

401 return stream.collect(count, include_meta=include_meta, sources=sources)