from __future__ import annotations import threading import time import unittest from app.state import RuntimeStats from app.tzsp_rust import RustTZSPReceiver class RustTZSPReceiverTelemetryTests(unittest.TestCase): def _receiver(self, sink=None): return RustTZSPReceiver( binary="/does/not/start-in-unit-tests", telemetry_socket="/tmp/unused-tzsp-test.sock", stats=RuntimeStats(), stop_event=threading.Event(), throughput_sink=sink, ) def test_telemetry_sample_drives_live_rate_without_packet_bytes_in_python(self): written = [] receiver = self._receiver(written.append) now_ms = int(time.time() * 1000) receiver._ingest( { "type": "tzsp_sample", "engine": "rust", "ready": True, "ts_ms": now_ms, "interval_ms": 1000, "last_packet_ms": now_ms, "bytes_total": 125_000_000, "bytes_in": 100_000_000, "bytes_out": 25_000_000, "packets_total": 80_000, "traffic_counters": {"bytes_total": 125_000_000, "packets_total": 80_000}, "kernel_udp_drops": 3, "kernel_udp_drops_interval": 0, "queue_dropped_datagrams": 0, "queue_drops_interval": 0, "truncated_datagrams": 0, "truncated_interval": 0, "queue_depth_batches": 1, "queue_capacity_batches": 24, "queue_capacity_bytes": 64 * 1024 * 1024, "capture_efficiency_pct": 100.0, "rx_datagrams_interval": 80_000, "rx_bytes_interval": 126_000_000, "rx_thread_alive": True, "worker_thread_alive": True, "rcvbuf_bytes": 425_984, "batch_size": 256, "datagram_bytes": 12_288, } ) current = receiver.current_throughput(3600) self.assertEqual(current["current_bps"], 1_000_000_000) self.assertEqual(current["current_in_bps"], 800_000_000) self.assertEqual(current["current_out_bps"], 200_000_000) self.assertEqual(current["current_pps"], 80_000) self.assertEqual(current["kernel_udp_drops"], 3) self.assertEqual(current["current_ingress_bps"], 1_008_000_000) self.assertEqual(current["capture_efficiency_pct"], 100.0) self.assertEqual(current["queue_fill_pct"], 4.2) self.assertEqual(current["loss_pps"], 0) self.assertTrue(current["rx_thread_alive"]) self.assertTrue(current["worker_thread_alive"]) self.assertEqual(current["receiver_engine"], "rust") self.assertEqual(len(written), 1) self.assertEqual(written[0]["bytes_total"], 125_000_000) def test_ingress_and_inspection_are_reported_separately_when_pipeline_is_behind(self): receiver = self._receiver() now_ms = int(time.time() * 1000) receiver._last = { "ts_ms": now_ms, "interval_ms": 1000, "bytes_total": 31_250_000, "packets_total": 24_000, "rx_bytes_interval": 125_000_000, "rx_datagrams_interval": 80_000, "kernel_udp_drops_interval": 10, "queue_drops_interval": 90, "truncated_interval": 0, "queue_depth_batches": 20, "queue_capacity_batches": 22, "capture_efficiency_pct": 99.875, } current = receiver.current_throughput(900) self.assertEqual(current["current_bps"], 250_000_000) self.assertEqual(current["current_ingress_bps"], 1_000_000_000) self.assertEqual(current["inspection_ratio_pct"], 25.0) self.assertEqual(current["queue_fill_pct"], 90.9) self.assertEqual(current["loss_pps"], 100.0) def test_stale_sample_reports_zero_current_rate(self): receiver = self._receiver() receiver._last = { "ts_ms": int(time.time() * 1000) - 10_000, "interval_ms": 1000, "bytes_total": 125_000_000, "packets_total": 80_000, } current = receiver.current_throughput(900) self.assertFalse(current["current_sample_fresh"]) self.assertEqual(current["current_bps"], 0) self.assertEqual(current["current_pps"], 0) if __name__ == "__main__": unittest.main()