297 lines
9.8 KiB
Python
297 lines
9.8 KiB
Python
|
|
#!/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()
|