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

1"""Python API for the SignalK v2 History API. 

2 

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""" 

12 

13import fnmatch 

14import logging 

15import re 

16from collections.abc import Sequence 

17from pathlib import Path 

18from typing import Any, Literal, Self 

19 

20import nanoarrow as na 

21import niquests 

22 

23from .._arrow import ArrowTable 

24from ..errors import SignalKError 

25from ..net import CACHE_DIR, normalise_host 

26from . import _results 

27from ._time import TimeRange 

28 

29logger = logging.getLogger("signalk_cli") 

30 

31HISTORY_BASE = "/signalk/v2/api/history" 

32 

33AGGREGATION_METHODS = ( 

34 "average", 

35 "min", 

36 "max", 

37 "first", 

38 "last", 

39 "mid", 

40 "middle_index", 

41 "sma", 

42 "ema", 

43) 

44 

45Shape = Literal["long", "wide"] 

46 

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"^[^.+(){}|^$\\]+$") 

50 

51 

52# --------------------------------------------------------------------------- 

53# Path specs and patterns 

54# --------------------------------------------------------------------------- 

55 

56 

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. 

64 

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`. 

69 

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 

87 

88 if any(":" in p for p in paths): 

89 return ",".join(paths), False 

90 

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 

98 

99 

100def _is_pattern(path: str) -> bool: 

101 return ":" not in path and any(c in path for c in _PATTERN_CHARS) 

102 

103 

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)) 

112 

113 

114def match_paths(patterns: Sequence[str], available: Sequence[str]) -> list[str]: 

115 """Expand path patterns against a list of known paths. 

116 

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. 

123 

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) 

136 

137 

138# --------------------------------------------------------------------------- 

139# Results 

140# --------------------------------------------------------------------------- 

141 

142 

143class HistoryResult: 

144 """A /values response, with conversions to tables and statistics. 

145 

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 """ 

151 

152 def __init__(self, payload: dict, *, wide: bool = False) -> None: 

153 self.payload = payload 

154 self.wide = wide 

155 

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 ) 

162 

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. 

165 

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 ) 

179 

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) 

183 

184 

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) 

199 

200 

201# --------------------------------------------------------------------------- 

202# Client 

203# --------------------------------------------------------------------------- 

204 

205 

206class HistoryClient: 

207 """Client for a SignalK server's v2 History API. 

208 

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. 

220 

221 Raises: 

222 SignalKError: From any method, when a request fails. 

223 """ 

224 

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() 

244 

245 def __enter__(self) -> Self: 

246 return self 

247 

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

249 self.close() 

250 

251 def close(self) -> None: 

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

253 if self._owns_session: 

254 self._session.close() 

255 

256 # -- low level ---------------------------------------------------------- 

257 

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. 

262 

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 

277 

278 def _get_json(self, endpoint: str, params: dict | None = None) -> Any: 

279 return self.fetch(endpoint, params).json() 

280 

281 # -- providers ---------------------------------------------------------- 

282 

283 def providers(self) -> dict[str, dict[str, Any]]: 

284 """Registered history providers, as `{id: {"isDefault": bool, ...}}`.""" 

285 return self._get_json("_providers") 

286 

287 def default_provider(self) -> str: 

288 """The id of the server's default history provider.""" 

289 return self._get_json("_providers/_default")["id"] 

290 

291 def _provider_cache_file(self) -> Path: 

292 return CACHE_DIR / f"{re.sub(r'[^\w.-]', '_', self.host)}.provider" 

293 

294 @property 

295 def provider(self) -> str | None: 

296 """The provider used for requests: the one given, else the server's default. 

297 

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 

323 

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. 

328 

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 

336 

337 # -- discovery ---------------------------------------------------------- 

338 

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))) 

342 

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))) 

346 

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. 

351 

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)) 

359 

360 # -- data --------------------------------------------------------------- 

361 

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. 

375 

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 

391 

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]. 

405 

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]). 

418 

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) 

433 

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. 

447 

448 Both shapes have one row per timestamp and path, with a UTC 

449 `timestamp` column and a `path` column: 

450 

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`, ...). 

456 

457 !!! warning 

458 Wide column names may change in a future release. 

459 

460 The other arguments are as for 

461 [`values()`][signalk_cli.history.api.HistoryClient.values]. 

462 

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) 

476 

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. 

486 

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 ) 

495 

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))