"""Wire Ghost telemetry: real, non-content network metrics from the app host. Hard privacy boundary (spec §3.3): only counters and timings are ever read — interface byte counters from /proc/net/dev, TCP connect latency, DNS resolution timing. Packet payloads are never inspected. """ import asyncio import socket import time from dataclasses import dataclass, field # Hosts we measure connect latency against (timing only, no data exchanged # beyond the handshake). REFERENCE_HOSTS = [("1.1.1.1", 53), ("10.30.20.107", 11434)] REFERENCE_DNS = ["example.com", "cloudflare.com"] @dataclass class TelemetrySample: jitter_bytes_per_s: float = 0.0 latency_variance_ms: float = 0.0 latency_mean_ms: float = 0.0 dns_ms: float = 0.0 extra: dict = field(default_factory=dict) def as_dict(self) -> dict: return { "jitter_bytes_per_s": self.jitter_bytes_per_s, "latency_variance_ms": self.latency_variance_ms, "latency_mean_ms": self.latency_mean_ms, "dns_ms": self.dns_ms, } def parse_proc_net_dev(text: str) -> dict[str, tuple[int, int]]: """Parse /proc/net/dev into {interface: (rx_bytes, tx_bytes)}.""" counters: dict[str, tuple[int, int]] = {} for line in text.splitlines()[2:]: if ":" not in line: continue iface, rest = line.split(":", 1) fields = rest.split() if len(fields) >= 9: counters[iface.strip()] = (int(fields[0]), int(fields[8])) return counters def _total_bytes(counters: dict[str, tuple[int, int]]) -> int: return sum(rx + tx for iface, (rx, tx) in counters.items() if iface != "lo") async def _read_counters() -> int: def _read() -> int: with open("/proc/net/dev") as f: return _total_bytes(parse_proc_net_dev(f.read())) return await asyncio.to_thread(_read) async def _tcp_rtt(host: str, port: int, timeout: float = 1.5) -> float | None: """Connect-then-close round trip time in ms. Nothing is sent or read.""" start = time.perf_counter() try: reader, writer = await asyncio.wait_for( asyncio.open_connection(host, port), timeout=timeout ) writer.close() try: await writer.wait_closed() except (ConnectionError, OSError): pass return (time.perf_counter() - start) * 1000 except (asyncio.TimeoutError, ConnectionError, OSError): return None async def _dns_timing(name: str, timeout: float = 1.5) -> float | None: start = time.perf_counter() try: await asyncio.wait_for(asyncio.to_thread(socket.getaddrinfo, name, None), timeout) return (time.perf_counter() - start) * 1000 except (asyncio.TimeoutError, socket.gaierror, OSError): return None async def sample_network(period_s: float = 1.0) -> TelemetrySample: """Take one telemetry sample: byte-counter jitter over `period_s`, plus latency variance and DNS hesitation.""" before = await _read_counters() await asyncio.sleep(period_s) after = await _read_counters() jitter = abs(after - before) / period_s rtts = [ rtt for rtt in await asyncio.gather(*(_tcp_rtt(h, p) for h, p in REFERENCE_HOSTS)) if rtt is not None ] dns_times = [ ms for ms in await asyncio.gather(*(_dns_timing(n) for n in REFERENCE_DNS)) if ms is not None ] sample = TelemetrySample(jitter_bytes_per_s=jitter) if rtts: mean = sum(rtts) / len(rtts) sample.latency_mean_ms = mean sample.latency_variance_ms = sum((r - mean) ** 2 for r in rtts) / len(rtts) if dns_times: sample.dns_ms = sum(dns_times) / len(dns_times) return sample def detect_wire_spike( history: list[float], current: float, *, min_samples: int = 6 ) -> float | None: """Return the spike ratio when `current` is anomalous vs the rolling baseline of jitter samples — the Wire Ghost noticing something move. Three guards so idle-hour noise doesn't cry ghost: enough history to have a baseline, an absolute floor (20 KB/s) so near-silent links stay silent, and a statistical (3σ) + relative (2.5× mean) threshold. """ if len(history) < min_samples: return None mean = sum(history) / len(history) if mean <= 0 or current <= 20_000: return None std = (sum((x - mean) ** 2 for x in history) / len(history)) ** 0.5 if current > mean + 3 * std and current > mean * 2.5: return current / mean return None