Files
box/bin/identity-variance-test.py
T

297 lines
9.8 KiB
Python
Raw Normal View History

#!/usr/bin/env python3
"""Identity variance testing framework.
Sends gentle, controlled probes from each NetVM node's Warp identity and
measures per-identity response behavior so that future rate limits can be
attributed (per-identity vs per-IP vs time-based vs random).
Design goals:
- Very gentle traffic: 2 probes per node per cycle, default 10-min cycle.
- Never logs secrets: only sha256 hashes of WireGuard private keys.
- JSONL results log: one line per probe, machine-readable.
Usage:
identity-variance-test.py probe # run one probe cycle over all nodes
identity-variance-test.py analyze [--since HOURS] [--log PATH]
Probes per node:
1. https://www.cloudflare.com/cdn-cgi/trace (identity + egress IP Cloudflare sees)
2. https://muse.ai/ (production-relevant landing page)
Exit codes: 0 ok, 1 partial (some nodes failed), 2 fatal.
"""
import argparse
import hashlib
import json
import os
import re
import subprocess
import sys
import time
from datetime import datetime, timezone
DEFAULT_LOG = os.path.expanduser(
"~/Projects/NetVM/.state/identity-variance/variance.jsonl"
)
NODES = ["muse", "pip", "646", "opm", "def"]
PROBES = [
("cf-trace", "https://www.cloudflare.com/cdn-cgi/trace"),
("muse-landing", "https://muse.ai/"),
]
CURL_TIMEOUT = 15
def sh(cmd, timeout=30):
"""Run cmd (list) and return (rc, stdout, stderr)."""
try:
p = subprocess.run(
cmd, capture_output=True, text=True, timeout=timeout
)
return p.returncode, p.stdout, p.stderr
except subprocess.TimeoutExpired as e:
return 124, (e.stdout or ""), "timeout"
except FileNotFoundError:
return 127, "", "command not found"
def node_netns_exists(node):
rc, out, _ = sh(["sudo", "-n", "ip", "netns", "list"])
return rc == 0 and f"warp-{node}" in out
def identity_hash(node):
"""Return truncated sha256 of the node's WireGuard private key (never the key)."""
rc, out, _ = sh(["sudo", "-n", "cat", f"/etc/netvm/{node}.conf"])
if rc == 0:
for line in out.splitlines():
m = re.match(r"\s*PrivateKey\s*=\s*(\S+)", line)
if m:
return hashlib.sha256(m.group(1).encode()).hexdigest()[:16]
return None
def probe_node(node, target_name, url):
"""One probe from inside the node's netns. Returns dict."""
# -w fields: http_code, time_total, size_download, remote_ip
fmt = "%{http_code} %{time_total} %{size_download} %{remote_ip}"
cmd = [
"sudo", "-n", "ip", "netns", "exec", f"warp-{node}",
"curl", "-s", "-o", "/dev/null", "-m", str(CURL_TIMEOUT),
"-w", fmt, url,
]
started = datetime.now(timezone.utc)
rc, out, err = sh(cmd, timeout=CURL_TIMEOUT + 10)
ended = datetime.now(timezone.utc)
rec = {
"ts": started.isoformat(),
"node": node,
"probe": target_name,
"url": url,
"curl_rc": rc,
}
parts = out.strip().split()
if rc == 0 and len(parts) == 4:
try:
rec["http_status"] = int(parts[0])
except ValueError:
rec["http_status"] = None
try:
rec["response_ms"] = round(float(parts[1]) * 1000, 1)
except ValueError:
rec["response_ms"] = None
try:
rec["bytes"] = int(parts[2])
except ValueError:
rec["bytes"] = None
rec["remote_ip"] = parts[3] if parts[3] != "0.0.0.0" else None
else:
rec["http_status"] = None
rec["response_ms"] = None
rec["bytes"] = None
rec["remote_ip"] = None
rec["curl_error"] = (err or "curl failed").strip()[:200]
rec["rate_limited"] = rec["http_status"] in (429,)
rec["blocked"] = rec["http_status"] in (403,)
rec["ok"] = rec["http_status"] is not None and 200 <= rec["http_status"] < 400
rec["elapsed_wall_ms"] = round(
(ended - started).total_seconds() * 1000, 1
)
return rec
def egress_ip(node):
"""Best-effort egress IP as seen from inside the netns."""
rc, out, _ = sh([
"sudo", "-n", "ip", "netns", "exec", f"warp-{node}",
"curl", "-s", "-m", "10", "https://api.ipify.org",
], timeout=20)
if rc == 0 and re.fullmatch(r"[0-9a-fA-F.:]+", out.strip()):
return out.strip()
return None
def run_probe_cycle(log_path):
os.makedirs(os.path.dirname(log_path), exist_ok=True)
results = []
any_ok, any_fail = False, False
idhash = {}
for node in NODES:
if not node_netns_exists(node):
results.append({
"ts": datetime.now(timezone.utc).isoformat(),
"node": node, "probe": "node-skip",
"ok": False, "note": "netns warp-%s missing" % node,
})
any_fail = True
continue
idhash[node] = identity_hash(node)
for target_name, url in PROBES:
rec = probe_node(node, target_name, url)
rec["identity_hash"] = idhash[node]
results.append(rec)
if rec["ok"]:
any_ok = True
else:
any_fail = True
time.sleep(1) # gentle pacing between probes
# Attach egress IP per node (one lookup per node, cached per cycle).
egress = {}
for node in {r["node"] for r in results if r.get("probe") != "node-skip"}:
egress[node] = egress_ip(node)
for rec in results:
if rec.get("probe") != "node-skip":
rec["egress_ip"] = egress.get(rec["node"])
with open(log_path, "a") as f:
for rec in results:
f.write(json.dumps(rec) + "\n")
print(json.dumps({
"cycle_ts": datetime.now(timezone.utc).isoformat(),
"log": log_path,
"records": len(results),
"ok": sum(1 for r in results if r.get("ok")),
"failed": sum(1 for r in results if not r.get("ok")),
"nodes": sorted({r["node"] for r in results}),
"egress": egress,
"identities": idhash,
}, indent=2))
if not any_ok:
return 2
return 1 if any_fail else 0
def analyze(log_path, since_hours=None):
if not os.path.exists(log_path):
print("no log yet at %s" % log_path, file=sys.stderr)
return 2
cutoff = None
if since_hours:
cutoff = time.time() - since_hours * 3600
per_node = {}
rate_limit_events = []
identity_changes = {}
egress_by_node = {}
with open(log_path) as f:
for line in f:
line = line.strip()
if not line:
continue
try:
r = json.loads(line)
except json.JSONDecodeError:
continue
if r.get("probe") == "node-skip":
continue
try:
ts = datetime.fromisoformat(r["ts"]).timestamp()
except (ValueError, KeyError):
continue
if cutoff and ts < cutoff:
continue
node = r["node"]
st = per_node.setdefault(node, {
"probes": 0, "ok": 0, "ms": [], "429": 0, "403": 0,
"statuses": {}, "identities": set(), "probes_by_target": {},
})
st["probes"] += 1
if r.get("ok"):
st["ok"] += 1
if r.get("response_ms") is not None:
st["ms"].append(r["response_ms"])
if r.get("rate_limited"):
st["429"] += 1
rate_limit_events.append((r["ts"], node, r["probe"]))
if r.get("blocked"):
st["403"] += 1
s = r.get("http_status")
st["statuses"][str(s)] = st["statuses"].get(str(s), 0) + 1
if r.get("identity_hash"):
st["identities"].add(r["identity_hash"])
t = r.get("probe")
st["probes_by_target"][t] = st["probes_by_target"].get(t, 0) + 1
if r.get("egress_ip"):
egress_by_node.setdefault(node, set()).add(r["egress_ip"])
def pct(vals, p):
if not vals:
return None
s = sorted(vals)
return round(s[min(len(s) - 1, int(p / 100 * len(s)))], 1)
report = {"log": log_path, "nodes": {}}
for node in sorted(per_node):
st = per_node[node]
report["nodes"][node] = {
"probes": st["probes"],
"ok_rate": round(st["ok"] / st["probes"], 3) if st["probes"] else 0,
"latency_ms": {
"mean": round(sum(st["ms"]) / len(st["ms"]), 1) if st["ms"] else None,
"p50": pct(st["ms"], 50),
"p99": pct(st["ms"], 99),
},
"http_429": st["429"],
"http_403": st["403"],
"statuses": st["statuses"],
"identity_rotations": max(0, len(st["identities"]) - 1),
"egress_ips": sorted(egress_by_node.get(node, set())),
}
if len(st["identities"]) > 1:
identity_changes[node] = sorted(st["identities"])
# Divergence analysis: did all nodes see the same fate at the same time?
report["rate_limit_events"] = [
{"ts": ts, "node": n, "probe": p} for ts, n, p in rate_limit_events[-50:]
]
report["identity_changes"] = identity_changes
print(json.dumps(report, indent=2))
return 0
def main():
ap = argparse.ArgumentParser(description="Identity variance testing framework")
sub = ap.add_subparsers(dest="cmd", required=True)
p_probe = sub.add_parser("probe", help="run one probe cycle")
p_probe.add_argument("--log", default=DEFAULT_LOG)
p_an = sub.add_parser("analyze", help="summarize variance data")
p_an.add_argument("--log", default=DEFAULT_LOG)
p_an.add_argument("--since", type=float, default=None,
help="only include last N hours")
args = ap.parse_args()
if args.cmd == "probe":
sys.exit(run_probe_cycle(args.log))
sys.exit(analyze(args.log, args.since))
if __name__ == "__main__":
main()