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

1"""CSV, JSON Lines, bare-value and Feather writers for the stream CLI.""" 

2 

3import csv 

4import json 

5from collections.abc import Sequence 

6from typing import IO 

7 

8from .._arrow import as_text 

9from .api import DeltaRow 

10 

11CSV_COLUMNS = ["timestamp", "context", "source", "path", "value"] 

12CSV_COLUMNS_WITH_KIND = ["timestamp", "context", "source", "path", "kind", "value"] 

13 

14 

15def columns(include_meta: bool) -> list[str]: 

16 return CSV_COLUMNS_WITH_KIND if include_meta else CSV_COLUMNS 

17 

18 

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 

25 

26 

27def write_csv_header(sink: IO[str], *, include_meta: bool = False) -> None: 

28 csv.writer(sink).writerow(columns(include_meta)) 

29 sink.flush() 

30 

31 

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) 

41 

42 

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) 

53 

54 

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)