113 lines
4.4 KiB
Python
113 lines
4.4 KiB
Python
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()
|