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
« 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.
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
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"""
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
23import niquests
25from .._arrow import ArrowTable
26from ..errors import SignalKError
27from ..net import normalise_host
29STREAM_PATH = "/signalk/v1/stream"
31SUBSCRIBE_POLICIES = ["none", "self", "all"]
32SUBSCRIPTION_POLICIES = ["instant", "ideal", "fixed"]
34Subscribe = Literal["none", "self", "all"]
35Policy = Literal["instant", "ideal", "fixed"]
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, "", "", ""))
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.
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").
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
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 }
85# ---------------------------------------------------------------------------
86# Messages and rows
87# ---------------------------------------------------------------------------
90def source_matches(source: str, patterns: Sequence[str]) -> bool:
91 """Match a `$source` string against source filter patterns (OR'd).
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 )
106def _update_source(update: dict) -> str:
107 return update.get("$source") or json.dumps(update.get("source", {}))
110class DeltaRow(NamedTuple):
111 """One value (or meta entry) from a delta message.
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 """
118 timestamp: str
119 context: str
120 source: str
121 path: str
122 kind: str
123 value: Any
126@dataclass(frozen=True)
127class DeltaMessage:
128 """A delta message from the server.
130 Attributes:
131 text: The message exactly as received.
132 payload: The decoded message.
133 """
135 text: str
136 payload: dict
138 def rows(
139 self, *, include_meta: bool = False, sources: Sequence[str] = ()
140 ) -> list[DeltaRow]:
141 """Flatten into rows, one per value.
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
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 )
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`.
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 )
202# ---------------------------------------------------------------------------
203# Stream and client
204# ---------------------------------------------------------------------------
207class DeltaStream:
208 """An open subscription. Iterate it for [`DeltaMessage`][signalk_cli.stream.api.DeltaMessage] objects.
210 Use as a context manager, or call [`close()`][signalk_cli.stream.api.DeltaStream.close], to close the connection.
212 Raises:
213 SignalKError: While iterating, if the connection is lost.
214 """
216 def __init__(self, ws: Any) -> None:
217 self._ws = ws
219 def __enter__(self) -> Self:
220 return self
222 def __exit__(self, *exc: object) -> None:
223 self.close()
225 def close(self) -> None:
226 """Close the WebSocket connection."""
227 self._ws.close()
229 def __iter__(self) -> Iterator[DeltaMessage]:
230 return self.messages()
232 def messages(self, count: int | None = None) -> Iterator[DeltaMessage]:
233 """Yield delta messages, skipping control messages such as the server's hello.
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
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)
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.
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)
283class StreamClient:
284 """Client for a SignalK server's v1 Streaming (delta) API.
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 """
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()
313 def __enter__(self) -> Self:
314 return self
316 def __exit__(self, *exc: object) -> None:
317 self.close()
319 def close(self) -> None:
320 """Close the session, if the client created it."""
321 if self._owns_session:
322 self._session.close()
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].
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).
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)
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.
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)