Coverage for src/signalk_cli/stream/cli.py: 99%

104 statements  

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

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

2 

3import re 

4import sys 

5from datetime import UTC, datetime 

6from pathlib import Path 

7from urllib.parse import urlparse 

8 

9import click 

10 

11from .._arrow import FEATHER_EXTENSIONS, write_feather 

12from .._cli import bare_option, host_option, resolve_host, stderr_ctx 

13from ..errors import SignalKError 

14from .api import ( 

15 SUBSCRIBE_POLICIES, 

16 SUBSCRIPTION_POLICIES, 

17 DeltaRow, 

18 StreamClient, 

19 rows_to_arrow, 

20) 

21from .output import ( 

22 write_csv_header, 

23 write_csv_rows, 

24 write_json_rows, 

25 write_values_rows, 

26) 

27 

28_AUTO_OUTPUT = "__auto_output__" 

29 

30# --------------------------------------------------------------------------- 

31# CLI group 

32# --------------------------------------------------------------------------- 

33 

34 

35@click.group(context_settings={"help_option_names": ["-h", "--help"]}) 

36def cli(): 

37 """SignalK v1 streaming (delta) CLI.""" 

38 

39 

40# --------------------------------------------------------------------------- 

41# deltas 

42# --------------------------------------------------------------------------- 

43 

44 

45@cli.command() 

46@click.argument("paths", nargs=-1, required=False, metavar="PATH...") 

47@host_option 

48@click.option("--no-cache", is_flag=True, help="Ignore cached host") 

49@click.option( 

50 "--context", 

51 "-c", 

52 default="vessels.self", 

53 show_default=True, 

54 help="SignalK context for the explicit subscription this command " 

55 "always sends. Also accepts the SignalK wildcard '*' (or 'vessels.*') " 

56 "to subscribe to every vessel at your own --policy/--period — the " 

57 "alternative to --subscribe all, which uses the server's default rate.", 

58) 

59@click.option( 

60 "--subscribe", 

61 type=click.Choice(SUBSCRIBE_POLICIES, case_sensitive=False), 

62 default="none", 

63 show_default=True, 

64 help="Connection-level auto-subscribe at the server's OWN default " 

65 "policy/period — separate from, and in addition to, this command's " 

66 "explicit --context subscription above. Defaults to 'none' (the " 

67 "SignalK spec's own default is 'self') to avoid double-subscribing " 

68 "your own context. 'all' adds other vessels at the server's rate; for " 

69 "other vessels at your chosen rate, use --context '*' instead.", 

70) 

71@click.option( 

72 "--policy", 

73 type=click.Choice(SUBSCRIPTION_POLICIES, case_sensitive=False), 

74 default="ideal", 

75 show_default=True, 

76 help="Per-path subscribe policy (SignalK Subscription Protocol 'policy' " 

77 "field): 'instant' sends every change, throttled by --min-period; " 

78 "'ideal' (default) behaves like instant but resends the last value if " 

79 "nothing changes within --period; 'fixed' always sends the last known " 

80 "value every --period regardless of changes.", 

81) 

82@click.option( 

83 "--period", 

84 type=float, 

85 default=60.0, 

86 show_default=True, 

87 metavar="SECONDS", 

88 help="Subscribe period in seconds — the resend interval used by the " 

89 "'ideal'/'fixed' policies. Converted to milliseconds on the wire.", 

90) 

91@click.option( 

92 "--min-period", 

93 "min_period", 

94 type=float, 

95 default=None, 

96 metavar="SECONDS", 

97 help="Fastest allowed transmission rate in seconds; only meaningful " 

98 "with --policy instant. Converted to milliseconds on the wire.", 

99) 

100@click.option( 

101 "--format", 

102 "fmt", 

103 default=None, 

104 type=click.Choice( 

105 ["csv", "json", "raw", "feather", "values"], case_sensitive=False 

106 ), 

107 help="Output format (default: inferred from --output extension, else csv). " 

108 "json is JSON Lines (one row object per line). raw is the exact delta " 

109 "message text, one per line. values is the bare value only, one per " 

110 "line — no timestamp/context/source/path/kind columns. feather requires " 

111 "pip install 'signalk-cli[feather]' and --output (cannot stream to stdout).", 

112) 

113@click.option("--no-header", is_flag=True, help="Suppress header row (CSV only)") 

114@click.option( 

115 "--include-meta", 

116 is_flag=True, 

117 help="Also emit rows for 'meta' entries (units, description, zones, etc.), " 

118 "not just 'values'. Adds a 'kind' column (value/meta) to csv/json/feather " 

119 "output. Ignored for --format raw, which always includes meta as-is.", 

120) 

121@click.option( 

122 "--source", 

123 "source", 

124 multiple=True, 

125 metavar="PATTERN", 

126 help="Only include updates whose $source matches PATTERN. Repeatable " 

127 "(OR'd together). PATTERN is a substring match unless it contains a " 

128 "glob metacharacter (*/?/[), in which case it's matched as a glob, " 

129 "e.g. --source Teltonika or --source '*.GP'. Filtering is client-side, " 

130 "applied after receipt — for --format raw (whole message, verbatim) a " 

131 "message passes if ANY of its updates match; other formats filter " 

132 "per-update.", 

133) 

134@click.option( 

135 "--output", 

136 "-o", 

137 is_flag=False, 

138 flag_value=_AUTO_OUTPUT, 

139 default=None, 

140 metavar="FILE", 

141 help="Write to FILE. Omit for stdout (default). Give without a filename to " 

142 "auto-name the file. Required for --format feather.", 

143) 

144@click.option( 

145 "--follow", 

146 "-f", 

147 is_flag=True, 

148 help="Keep streaming until interrupted (Ctrl-C) or --count is reached. " 

149 "Without this, print the next message then exit.", 

150) 

151@click.option( 

152 "--count", 

153 "-n", 

154 type=int, 

155 default=None, 

156 metavar="N", 

157 help="Number of delta messages to output. Default: 1 without --follow, " 

158 "unlimited with --follow.", 

159) 

160@bare_option 

161def deltas( 

162 paths, 

163 host, 

164 no_cache, 

165 context, 

166 subscribe, 

167 policy, 

168 period, 

169 min_period, 

170 fmt, 

171 no_header, 

172 include_meta, 

173 source, 

174 output, 

175 follow, 

176 count, 

177 bare, 

178): 

179 """Stream live delta updates from the SignalK v1 Streaming API. 

180 

181 Connects via WebSocket and prints delta messages as they arrive, in 

182 csv, json (JSON Lines), raw, values (bare value only), or Feather 

183 format. 

184 

185 Always sends an explicit subscribe message for --context, covering 

186 PATH arguments if given, otherwise every path ('*'). PATH arguments 

187 are sent verbatim to the server, one per path. They may be literal 

188 SignalK paths (e.g. navigation.speedOverGround) or contain the 

189 SignalK subscription wildcard '*', matched server-side per the 

190 Subscription Protocol: '*' at the end of a path matches any suffix 

191 (navigation.*), and '*' as a middle segment matches any single 

192 segment there (propulsion.*.oilTemperature). Quote wildcarded paths 

193 to stop the shell expanding them. --policy/--period/--min-period 

194 control how that subscription behaves; --subscribe only controls 

195 whether the connection *additionally* auto-subscribes at the 

196 server's own default policy. 

197 

198 \b 

199 Examples: 

200 # Next update for one path, then exit 

201 signalk_cli.stream deltas --host 10.36.10.21 navigation.speedOverGround 

202 

203 # Tail all navigation updates every 5s until Ctrl-C 

204 signalk_cli.stream deltas --host 10.36.10.21 --follow --period 5 'navigation.*' 

205 

206 # Tail oil temperature across every engine, sent instantly on change 

207 signalk_cli.stream deltas --host 10.36.10.21 --follow --policy instant \\ 

208 'propulsion.*.oilTemperature' 

209 

210 # Next 20 messages across all paths, as JSON Lines 

211 signalk_cli.stream deltas --host 10.36.10.21 --format json --count 20 

212 

213 # Capture 100 messages to a Feather file (requires signalk-cli[feather]) 

214 signalk_cli.stream deltas --host 10.36.10.21 --count 100 -o capture.feather 

215 

216 # Bare speed values from one sensor, piped straight into another tool 

217 signalk_cli.stream deltas --host 10.36.10.21 --follow --format values \\ 

218 --source Teltonika --bare navigation.speedOverGround 

219 """ 

220 with stderr_ctx(bare): 

221 client = StreamClient( 

222 resolve_host(host, no_cache), context=context, subscribe=subscribe 

223 ) 

224 

225 auto_name = output == _AUTO_OUTPUT 

226 

227 if fmt is None: 

228 if output and not auto_name and output != "-": 

229 suffix = Path(output).suffix.lower() 

230 if suffix in FEATHER_EXTENSIONS: 

231 fmt = "feather" 

232 elif suffix == ".json": 

233 fmt = "json" 

234 else: 

235 fmt = "csv" 

236 else: 

237 fmt = "csv" 

238 

239 if auto_name: 

240 server_name = urlparse(client.host).hostname or re.sub( 

241 r"[^\w.-]", "_", client.host 

242 ) 

243 ts = datetime.now(UTC).strftime("%Y%m%dT%H%M%SZ") 

244 ext = ( 

245 ".feather" 

246 if fmt == "feather" 

247 else ".json" 

248 if fmt in ("json", "raw") 

249 else ".txt" 

250 if fmt == "values" 

251 else ".csv" 

252 ) 

253 output = f"signalk-stream-{server_name}-{ts}{ext}" 

254 

255 write_to_stdout = output is None or output == "-" 

256 write_to_file = not write_to_stdout 

257 

258 if fmt == "feather" and write_to_stdout: 

259 raise click.UsageError( 

260 "feather cannot be written to stdout (binary format); " 

261 "use --output FILE or --output to auto-name" 

262 ) 

263 

264 click.echo(f"Server: {client.host}", err=True) 

265 click.echo(f"Context: {context}", err=True) 

266 click.echo(f"Subscribe: {subscribe}", err=True) 

267 min_period_note = ( 

268 f", min_period={min_period}s" if min_period is not None else "" 

269 ) 

270 click.echo( 

271 f"Policy: {policy} (period={period}s{min_period_note})", err=True 

272 ) 

273 click.echo(f"Format: {fmt}", err=True) 

274 

275 try: 

276 stream = client.open( 

277 paths, 

278 policy=policy, 

279 period=period, 

280 min_period=min_period, 

281 timeout=None if follow else 30, 

282 ) 

283 except SignalKError as e: 

284 click.echo(f"Error connecting to stream: {e}", err=True) 

285 sys.exit(1) 

286 

287 effective_count = count if count is not None else (None if follow else 1) 

288 

289 message_count = 0 

290 row_total = 0 

291 header_written = False 

292 feather_rows: list[DeltaRow] = [] 

293 fh = ( 

294 open(output, "w", newline="") # noqa: SIM115 

295 if write_to_file and fmt != "feather" 

296 else None 

297 ) 

298 sink = fh or sys.stdout 

299 try: 

300 for message in stream.messages(effective_count): 

301 message_count += 1 

302 if fmt == "raw": 

303 if message.matches_sources(source): 

304 click.echo(message.text, file=sink) 

305 continue 

306 rows = message.rows(include_meta=include_meta, sources=source) 

307 if fmt == "feather": 

308 feather_rows.extend(rows) 

309 row_total = len(feather_rows) 

310 elif fmt == "json": 

311 row_total += write_json_rows(rows, sink, include_meta=include_meta) 

312 elif fmt == "values": 

313 row_total += write_values_rows(rows, sink) 

314 else: 

315 if not header_written and not no_header: 

316 write_csv_header(sink, include_meta=include_meta) 

317 header_written = True 

318 row_total += write_csv_rows(rows, sink, include_meta=include_meta) 

319 except KeyboardInterrupt: 

320 pass 

321 except SignalKError as e: 

322 click.echo(f"Stream connection lost: {e}", err=True) 

323 finally: 

324 stream.close() 

325 if fh: 

326 fh.close() 

327 

328 if fmt == "feather": 

329 write_feather( 

330 rows_to_arrow(feather_rows, include_meta=include_meta), output 

331 ) 

332 

333 if write_to_file: 

334 click.echo(f"Wrote {output}", err=True) 

335 

336 if fmt == "raw": 

337 click.echo(f"{message_count} message(s)", err=True) 

338 else: 

339 click.echo(f"{message_count} message(s), {row_total} row(s)", err=True)