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
« 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."""
3import re
4import sys
5from datetime import UTC, datetime
6from pathlib import Path
7from urllib.parse import urlparse
9import click
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)
28_AUTO_OUTPUT = "__auto_output__"
30# ---------------------------------------------------------------------------
31# CLI group
32# ---------------------------------------------------------------------------
35@click.group(context_settings={"help_option_names": ["-h", "--help"]})
36def cli():
37 """SignalK v1 streaming (delta) CLI."""
40# ---------------------------------------------------------------------------
41# deltas
42# ---------------------------------------------------------------------------
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.
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.
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.
198 \b
199 Examples:
200 # Next update for one path, then exit
201 signalk_cli.stream deltas --host 10.36.10.21 navigation.speedOverGround
203 # Tail all navigation updates every 5s until Ctrl-C
204 signalk_cli.stream deltas --host 10.36.10.21 --follow --period 5 'navigation.*'
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'
210 # Next 20 messages across all paths, as JSON Lines
211 signalk_cli.stream deltas --host 10.36.10.21 --format json --count 20
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
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 )
225 auto_name = output == _AUTO_OUTPUT
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"
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}"
255 write_to_stdout = output is None or output == "-"
256 write_to_file = not write_to_stdout
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 )
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)
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)
287 effective_count = count if count is not None else (None if follow else 1)
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()
328 if fmt == "feather":
329 write_feather(
330 rows_to_arrow(feather_rows, include_meta=include_meta), output
331 )
333 if write_to_file:
334 click.echo(f"Wrote {output}", err=True)
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)