Coverage for src/signalk_cli/history/api.py: 99%
178 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 v2 History API.
3Examples:
4 ```python
5 from signalk_cli.history import HistoryClient, TimeRange
6 with HistoryClient("http://boat.local:3000") as client:
7 table = client.query(["navigation.speedOverGround"], TimeRange(duration="PT1H"))
8 import polars as pl
9 df = pl.DataFrame(table)
10 ```
11"""
13import fnmatch
14import logging
15import re
16from collections.abc import Sequence
17from pathlib import Path
18from typing import Any, Literal, Self
20import nanoarrow as na
21import niquests
23from .._arrow import ArrowTable
24from ..errors import SignalKError
25from ..net import CACHE_DIR, normalise_host
26from . import _results
27from ._time import TimeRange
29logger = logging.getLogger("signalk_cli")
31HISTORY_BASE = "/signalk/v2/api/history"
33AGGREGATION_METHODS = (
34 "average",
35 "min",
36 "max",
37 "first",
38 "last",
39 "mid",
40 "middle_index",
41 "sma",
42 "ema",
43)
45Shape = Literal["long", "wide"]
47# "." is left out: it's in every SignalK path, so a plain dotted path is a literal
48_PATTERN_CHARS = set(r"*+?[](){}|^$\\")
49_GLOB_ONLY_RE = re.compile(r"^[^.+(){}|^$\\]+$")
52# ---------------------------------------------------------------------------
53# Path specs and patterns
54# ---------------------------------------------------------------------------
57def build_path_specs(
58 paths: Sequence[str],
59 aggregation: str | None = None,
60 samples: int | None = None,
61 alpha: float | None = None,
62) -> tuple[str, bool]:
63 """Build the `paths` query parameter, with an aggregation method per path.
65 Paths that already carry a method (`navigation.speedOverGround:sma:5`)
66 pass through unchanged. With no aggregation and no inline methods, each
67 path is requested as min, average and max (wide shape), except position
68 paths, which don't support those and are requested with `mid`.
70 Returns:
71 `(paths_param, wide)`, where `wide` says the response should be
72 read in the wide shape.
73 """
74 if aggregation:
75 specs = []
76 for path in paths:
77 if ":" in path:
78 specs.append(path)
79 continue
80 spec = f"{path}:{aggregation}"
81 if aggregation == "sma" and samples is not None:
82 spec += f":{samples}"
83 elif aggregation == "ema" and alpha is not None:
84 spec += f":{alpha}"
85 specs.append(spec)
86 return ",".join(specs), False
88 if any(":" in p for p in paths):
89 return ",".join(paths), False
91 specs = []
92 for p in paths:
93 if _results.POSITION_RE.fullmatch(p):
94 specs.append(f"{p}:mid")
95 else:
96 specs.extend(f"{p}:{m}" for m in ("min", "average", "max"))
97 return ",".join(specs), True
100def _is_pattern(path: str) -> bool:
101 return ":" not in path and any(c in path for c in _PATTERN_CHARS)
104def _compile(pattern: str) -> re.Pattern:
105 if _GLOB_ONLY_RE.match(pattern):
106 return re.compile(fnmatch.translate(pattern))
107 try:
108 return re.compile(pattern)
109 except re.PatternError:
110 logger.info("Note: '%s' is not valid regex, treating as glob", pattern)
111 return re.compile(fnmatch.translate(pattern))
114def match_paths(patterns: Sequence[str], available: Sequence[str]) -> list[str]:
115 """Expand path patterns against a list of known paths.
117 Literal paths and inline specs (`path:method`) pass through unchanged.
118 A path is a pattern if it contains any of `* ? + [ ] ( ) { } | ^ $ \\`
119 (a `.` alone doesn't count, so `navigation.speedOverGround` is literal).
120 Patterns are matched as globs (`navigation.*`) when they only use glob
121 characters, otherwise as Python regular expressions (searched, not
122 anchored); an invalid regex falls back to a glob.
124 Returns:
125 The literal paths in the order given, followed by all matched paths, sorted.
126 """
127 literals = [p for p in patterns if not _is_pattern(p)]
128 matched: set[str] = set()
129 for pattern in (p for p in patterns if _is_pattern(p)):
130 rx = _compile(pattern)
131 hits = {path for path in available if rx.search(path)}
132 if not hits:
133 logger.warning("Warning: '%s' matched no paths", pattern)
134 matched |= hits
135 return literals + sorted(matched)
138# ---------------------------------------------------------------------------
139# Results
140# ---------------------------------------------------------------------------
143class HistoryResult:
144 """A /values response, with conversions to tables and statistics.
146 Attributes:
147 payload: The decoded JSON response, as the server sent it.
148 wide: True when the request asked for min/average/max per path, so
149 the result reads naturally in the wide shape.
150 """
152 def __init__(self, payload: dict, *, wide: bool = False) -> None:
153 self.payload = payload
154 self.wide = wide
156 @property
157 def paths(self) -> list[str]:
158 """Distinct paths in the response, in the order the server listed them."""
159 return list(
160 dict.fromkeys(v.get("path", "") for v in self.payload.get("values", []))
161 )
163 def to_arrow(self, shape: Shape | None = None) -> ArrowTable:
164 """Convert to a table; see [`query()`][signalk_cli.history.api.HistoryClient.query] for the two shapes.
166 Args:
167 shape: `"long"` or `"wide"`. Defaults to wide when the request
168 asked for min/average/max, otherwise long.
169 """
170 if (shape or ("wide" if self.wide else "long")) == "wide":
171 timestamps, paths, cols = _results.wide_rows(self.payload)
172 return ArrowTable.from_rows(
173 timestamps, {"path": paths, **cols}, text_columns=["path"]
174 )
175 timestamps, paths, values = _results.long_rows(self.payload)
176 return ArrowTable.from_rows(
177 timestamps, {"path": paths, "value": values}, text_columns=["path"]
178 )
180 def cardinality(self) -> list[dict[str, Any]]:
181 """Per-path statistics; see [`cardinality()`][signalk_cli.history.api.HistoryClient.cardinality]."""
182 return _results.cardinality(self.payload)
185def _cardinality_table(rows: list[dict[str, Any]]) -> ArrowTable:
186 columns: dict[str, tuple[list, Any]] = {
187 "path": ([r["path"] for r in rows], na.string())
188 }
189 for col in _results.CARDINALITY_COLUMNS[1:]:
190 values = [r[col] for r in rows]
191 if col in ("min", "max", "average"):
192 columns[col] = (
193 [None if v is None else float(v) for v in values],
194 na.float64(),
195 )
196 else:
197 columns[col] = (values, na.int64())
198 return ArrowTable(columns)
201# ---------------------------------------------------------------------------
202# Client
203# ---------------------------------------------------------------------------
206class HistoryClient:
207 """Client for a SignalK server's v2 History API.
209 Args:
210 host: Server URL, e.g. `http://boat.local:3000` (`http://` is
211 added if there's no scheme).
212 provider: History provider plugin id. Defaults to the server's
213 default provider, looked up on first use.
214 context: SignalK context to query, e.g. `vessels.self`.
215 session: A niquests session to send requests with. One is created
216 (and closed with the client) if not given.
217 cache: Remember each server's default provider on disk, under
218 `~/.cache/signalk-cli`, to save a request next time.
219 timeout: Seconds to wait for each response.
221 Raises:
222 SignalKError: From any method, when a request fails.
223 """
225 def __init__(
226 self,
227 host: str,
228 *,
229 provider: str | None = None,
230 context: str = "vessels.self",
231 session: niquests.Session | None = None,
232 cache: bool = False,
233 timeout: float = 60,
234 ) -> None:
235 self.host = normalise_host(host).rstrip("/")
236 self.base_url = self.host + HISTORY_BASE
237 self.context = context
238 self.timeout = timeout
239 self._provider = provider
240 self._provider_resolved = provider is not None
241 self._cache = cache
242 self._owns_session = session is None
243 self._session = session if session is not None else niquests.Session()
245 def __enter__(self) -> Self:
246 return self
248 def __exit__(self, *exc: object) -> None:
249 self.close()
251 def close(self) -> None:
252 """Close the session, if the client created it."""
253 if self._owns_session:
254 self._session.close()
256 # -- low level ----------------------------------------------------------
258 def fetch(
259 self, endpoint: str, params: dict | None = None, *, stream: bool = False
260 ) -> niquests.Response:
261 """GET a History API endpoint (e.g. `"values"`, `"paths"`) and return the response.
263 A low-level escape hatch for callers that want the raw response body;
264 the other methods are usually more convenient.
265 """
266 try:
267 resp = self._session.get(
268 f"{self.base_url}/{endpoint}",
269 params=params,
270 timeout=self.timeout,
271 stream=stream,
272 )
273 resp.raise_for_status()
274 except niquests.RequestException as e:
275 raise SignalKError.from_request(e) from e
276 return resp
278 def _get_json(self, endpoint: str, params: dict | None = None) -> Any:
279 return self.fetch(endpoint, params).json()
281 # -- providers ----------------------------------------------------------
283 def providers(self) -> dict[str, dict[str, Any]]:
284 """Registered history providers, as `{id: {"isDefault": bool, ...}}`."""
285 return self._get_json("_providers")
287 def default_provider(self) -> str:
288 """The id of the server's default history provider."""
289 return self._get_json("_providers/_default")["id"]
291 def _provider_cache_file(self) -> Path:
292 return CACHE_DIR / f"{re.sub(r'[^\w.-]', '_', self.host)}.provider"
294 @property
295 def provider(self) -> str | None:
296 """The provider used for requests: the one given, else the server's default.
298 The default is fetched once (or read from the disk cache if enabled).
299 If it can't be fetched, requests go without one and the server picks.
300 """
301 if self._provider_resolved:
302 return self._provider
303 self._provider_resolved = True
304 cache_file = self._provider_cache_file()
305 if self._cache:
306 try:
307 self._provider = cache_file.read_text().strip() or None
308 except OSError:
309 pass
310 if not self._provider:
311 try:
312 self._provider = self.default_provider()
313 except SignalKError as e:
314 logger.warning("Warning: could not fetch default provider: %s", e)
315 return None
316 if self._cache:
317 try:
318 CACHE_DIR.mkdir(parents=True, exist_ok=True)
319 cache_file.write_text(self._provider)
320 except OSError:
321 pass
322 return self._provider
324 def request_params(
325 self, time: TimeRange | None = None, **extra: Any
326 ) -> dict[str, Any]:
327 """Query parameters for a request: the time range, provider, and any extras given.
329 Useful with [`fetch()`][signalk_cli.history.api.HistoryClient.fetch]; `None` extras are left out.
330 """
331 params: dict[str, Any] = (time or TimeRange()).resolved().params()
332 if self.provider:
333 params["provider"] = self.provider
334 params.update({k: v for k, v in extra.items() if v is not None})
335 return params
337 # -- discovery ----------------------------------------------------------
339 def contexts(self, time: TimeRange | None = None) -> list[str]:
340 """Contexts (vessels, aircraft, ...) with data in the time range, sorted."""
341 return sorted(self._get_json("contexts", self.request_params(time)))
343 def paths(self, time: TimeRange | None = None) -> list[str]:
344 """Paths with data in the time range, sorted."""
345 return sorted(self._get_json("paths", self.request_params(time)))
347 def expand_paths(
348 self, patterns: Sequence[str], time: TimeRange | None = None
349 ) -> list[str]:
350 """Expand glob/regex patterns to the matching paths with data in the time range.
352 Only fetches the server's path list if there's a pattern to match.
353 See [`match_paths()`][signalk_cli.history.api.match_paths] for the matching rules.
354 """
355 if not any(_is_pattern(p) for p in patterns):
356 return list(patterns)
357 logger.info("Resolving patterns against server paths...")
358 return match_paths(patterns, self.paths(time))
360 # -- data ---------------------------------------------------------------
362 def value_params(
363 self,
364 paths: Sequence[str],
365 time: TimeRange | None = None,
366 *,
367 aggregation: str | None = None,
368 samples: int | None = None,
369 alpha: float | None = None,
370 resolution: str | int | None = None,
371 context: str | None = None,
372 expand: bool = True,
373 ) -> tuple[dict[str, Any], bool]:
374 """The query parameters [`values()`][signalk_cli.history.api.HistoryClient.values] sends, and whether the result is wide.
376 Useful with [`fetch()`][signalk_cli.history.api.HistoryClient.fetch] to get the raw response body. Arguments are
377 as for [`values()`][signalk_cli.history.api.HistoryClient.values].
378 """
379 time = (time or TimeRange()).resolved()
380 resolved = self.expand_paths(paths, time) if expand else list(paths)
381 if not resolved:
382 raise ValueError("No paths to query")
383 spec, wide = build_path_specs(resolved, aggregation, samples, alpha)
384 params = self.request_params(
385 time,
386 paths=spec,
387 context=context or self.context,
388 resolution=resolution,
389 )
390 return params, wide
392 def values(
393 self,
394 paths: Sequence[str],
395 time: TimeRange | None = None,
396 *,
397 aggregation: str | None = None,
398 samples: int | None = None,
399 alpha: float | None = None,
400 resolution: str | int | None = None,
401 context: str | None = None,
402 expand: bool = True,
403 ) -> HistoryResult:
404 """Fetch values for the given paths, as a [`HistoryResult`][signalk_cli.history.api.HistoryResult].
406 Args:
407 paths: Paths, glob/regex patterns, or inline specs such as
408 `navigation.speedOverGround:sma:5`.
409 time: Time range; defaults to the last hour.
410 aggregation: Method applied to each path without an inline spec,
411 one of [`AGGREGATION_METHODS`][signalk_cli.history.api.AGGREGATION_METHODS]. If neither this nor
412 inline specs are given, min/average/max are fetched (wide).
413 samples: Window size for `sma`.
414 alpha: Smoothing factor for `ema`.
415 resolution: Sample window, as seconds or an expression like `1m`.
416 context: Overrides the client's context for this request.
417 expand: Expand patterns in `paths` first (see [`expand_paths()`][signalk_cli.history.api.HistoryClient.expand_paths]).
419 Raises:
420 ValueError: If no paths are left to query after expansion.
421 """
422 params, wide = self.value_params(
423 paths,
424 time,
425 aggregation=aggregation,
426 samples=samples,
427 alpha=alpha,
428 resolution=resolution,
429 context=context,
430 expand=expand,
431 )
432 return HistoryResult(self._get_json("values", params), wide=wide)
434 def query(
435 self,
436 paths: Sequence[str],
437 time: TimeRange | None = None,
438 *,
439 shape: Shape | None = None,
440 aggregation: str | None = None,
441 samples: int | None = None,
442 alpha: float | None = None,
443 resolution: str | int | None = None,
444 context: str | None = None,
445 ) -> ArrowTable:
446 """Fetch values as a table for polars, pandas, pyarrow, DuckDB, etc.
448 Both shapes have one row per timestamp and path, with a UTC
449 `timestamp` column and a `path` column:
451 - **long**: a single `value` column. It's float64 if every value is
452 a number, otherwise text, with objects and arrays as JSON.
453 - **wide**: `min_value`/`avg_value`/`max_value` for number paths,
454 and one column per element for array paths (`longitude`/`latitude`
455 for positions, otherwise `value_0`, `value_1`, ...).
457 !!! warning
458 Wide column names may change in a future release.
460 The other arguments are as for
461 [`values()`][signalk_cli.history.api.HistoryClient.values].
463 Args:
464 shape: Defaults to wide if no aggregation or inline spec is
465 given (min/average/max are fetched), otherwise long.
466 """
467 return self.values(
468 paths,
469 time,
470 aggregation=aggregation,
471 samples=samples,
472 alpha=alpha,
473 resolution=resolution,
474 context=context,
475 ).to_arrow(shape)
477 def cardinality(
478 self,
479 paths: Sequence[str] = ("*",),
480 time: TimeRange | None = None,
481 *,
482 resolution: str | int | None = None,
483 context: str | None = None,
484 ) -> ArrowTable:
485 """Per-path statistics over the time range, as a table.
487 Columns: `path`, `distinct_values`,
488 `distinct_values_2_decimal_places`, `nulls`, `zeroes`, `min`,
489 `max`, `average`. `min`/`max`/`average` are null unless every
490 value of the path is a number.
491 """
492 return _cardinality_table(
493 self.cardinality_rows(paths, time, resolution=resolution, context=context)
494 )
496 def cardinality_rows(
497 self,
498 paths: Sequence[str] = ("*",),
499 time: TimeRange | None = None,
500 *,
501 resolution: str | int | None = None,
502 context: str | None = None,
503 ) -> list[dict[str, Any]]:
504 """Like [`cardinality()`][signalk_cli.history.api.HistoryClient.cardinality], as a list of dicts."""
505 time = (time or TimeRange()).resolved()
506 resolved = self.expand_paths(paths, time)
507 if not resolved:
508 raise ValueError("No paths to query")
509 params = self.request_params(
510 time,
511 paths=",".join(resolved),
512 context=context or self.context,
513 resolution=resolution,
514 )
515 return _results.cardinality(self._get_json("values", params))