Coverage for src/signalk_cli/stream/output.py: 100%
38 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"""CSV, JSON Lines, bare-value and Feather writers for the stream CLI."""
3import csv
4import json
5from collections.abc import Sequence
6from typing import IO
8from .._arrow import as_text
9from .api import DeltaRow
11CSV_COLUMNS = ["timestamp", "context", "source", "path", "value"]
12CSV_COLUMNS_WITH_KIND = ["timestamp", "context", "source", "path", "kind", "value"]
15def columns(include_meta: bool) -> list[str]:
16 return CSV_COLUMNS_WITH_KIND if include_meta else CSV_COLUMNS
19def _text_row(row: DeltaRow, include_meta: bool) -> list[str]:
20 values = [row.timestamp, row.context, row.source, row.path]
21 if include_meta:
22 values.append(row.kind)
23 values.append(as_text(row.value) or "")
24 return values
27def write_csv_header(sink: IO[str], *, include_meta: bool = False) -> None:
28 csv.writer(sink).writerow(columns(include_meta))
29 sink.flush()
32def write_csv_rows(
33 rows: Sequence[DeltaRow], sink: IO[str], *, include_meta: bool = False
34) -> int:
35 """Write rows as CSV lines. Returns the number of rows written."""
36 writer = csv.writer(sink)
37 for row in rows:
38 writer.writerow(_text_row(row, include_meta))
39 sink.flush()
40 return len(rows)
43def write_json_rows(
44 rows: Sequence[DeltaRow], sink: IO[str], *, include_meta: bool = False
45) -> int:
46 """Write rows as JSON Lines (one row object per line). Returns row count."""
47 cols = columns(include_meta)
48 for row in rows:
49 sink.write(json.dumps(dict(zip(cols, _text_row(row, include_meta)))))
50 sink.write("\n")
51 sink.flush()
52 return len(rows)
55def write_values_rows(rows: Sequence[DeltaRow], sink: IO[str]) -> int:
56 """Write bare values, one per line — for piping one path's readings elsewhere."""
57 for row in rows:
58 sink.write(as_text(row.value) or "")
59 sink.write("\n")
60 sink.flush()
61 return len(rows)