/
/
/
1"""VBAN stream stats reporting."""
2
3from __future__ import annotations
4
5import asyncio
6import logging
7from time import monotonic
8
9
10class VBANStatsReporter:
11 """VBAN stream stats reporter."""
12
13 def __init__(self, pcm_sample_size: int, report_interval: int = 30) -> None:
14 """Initialize stats reporter."""
15 self._pcm_sample_size = pcm_sample_size
16 self._report_interval = report_interval
17 self._handle: asyncio.TimerHandle | None = None
18 self._loop = asyncio.get_running_loop()
19 self._logger = logging.getLogger(__name__)
20 self._active = False
21 self._stream_id: str | None = None
22 self._stream_details = ""
23 self._last_report_time = 0.0
24 self._packet_count = 0
25 self._kbytes_count = 0.0
26
27 def start(self, stream_id: str, stream_details: str) -> None:
28 """Start stats."""
29 self.cancel()
30 self._active = True
31 self._stream_id = stream_id
32 self._stream_details = stream_details
33 self._reset(stream_id, call_interval=5)
34
35 def _reset(self, instance_id: str, call_interval: int | None = None) -> None:
36 """Reset stats."""
37 if not self._active or instance_id != self._stream_id:
38 return
39 self._last_report_time = monotonic()
40 self._packet_count = 0
41 self._kbytes_count = 0.0
42 _interval = (
43 call_interval
44 if (call_interval and call_interval < self._report_interval)
45 else self._report_interval
46 )
47 self._handle = self._loop.call_later(_interval, self._report_cb, instance_id)
48
49 def update(self, instance_id: str, vban_bytes_len: int) -> None:
50 """Update stats."""
51 if not self._active or instance_id != self._stream_id:
52 return
53 self._packet_count += 1
54 self._kbytes_count += vban_bytes_len / 1024
55
56 def _report_cb(self, instance_id: str) -> None:
57 """Report stats callback."""
58 if not self._active or instance_id != self._stream_id:
59 return
60 try:
61 _secs_since_last_report = monotonic() - self._last_report_time
62 _avg_pkts = self._packet_count / _secs_since_last_report
63 _avg_kbytes = self._kbytes_count / _secs_since_last_report
64 _avg_pkt_length = (
65 (self._kbytes_count / self._packet_count) * 1024 if self._packet_count else 0
66 )
67 _target_kbytes = self._pcm_sample_size / 1024
68 _target_pkts = self._pcm_sample_size / _avg_pkt_length if _avg_pkt_length else 0
69 except ZeroDivisionError:
70 self._logger.exception("")
71 else:
72 self._logger.debug(
73 "VBAN stream stats: avg pkts/s processed/target: %d/%d // avg kB/s processed/target: %d/%d // avg pkt length: %d bytes // Stream %s",
74 _avg_pkts,
75 _target_pkts,
76 _avg_kbytes,
77 _target_kbytes,
78 _avg_pkt_length,
79 self._stream_details,
80 )
81 finally:
82 if self._active and instance_id == self._stream_id:
83 self._reset(instance_id)
84
85 def cancel(self, instance_id: str | None = None) -> None:
86 """Cancel stats reporter."""
87 if instance_id is not None and instance_id != self._stream_id:
88 return
89 self._active = False
90 self._stream_id = None
91 if self._handle:
92 self._handle.cancel()
93 self._handle = None
94