From 41499e069d14d93bebf7eb58bfa0f7071c6c8e26 Mon Sep 17 00:00:00 2001 From: operator Date: Thu, 8 Oct 2026 04:03:45 +0000 Subject: [PATCH] feat(identity): per-scope network identity plane (slices 1-5) Fingerprint map + pure resolver (account umbrella / key-level scope rule), live Warp provider on warp-* structures, broker lifecycle (up/down/cycle/exec/routes/status/bind), wireguard+socks boilerplate stubs, agent-manager bind integration. CLI carries emails and fingerprints only; key bytes never appear. 41 committed tests. --- .gitignore | 1 + bin/identity-broker.py | 342 +++++++++++++++++++++++++ bin/identity-provider.py | 237 +++++++++++++++++ bin/identity-resolve.py | 186 ++++++++++++++ docs/IDENTITY-PLANE.md | 65 +++++ identity-map.json | 22 ++ tests/test_identity_plane.py | 475 +++++++++++++++++++++++++++++++++++ 7 files changed, 1328 insertions(+) create mode 100755 bin/identity-broker.py create mode 100755 bin/identity-provider.py create mode 100755 bin/identity-resolve.py create mode 100644 docs/IDENTITY-PLANE.md create mode 100644 identity-map.json create mode 100644 tests/test_identity_plane.py diff --git a/.gitignore b/.gitignore index 3267519..9304318 100644 --- a/.gitignore +++ b/.gitignore @@ -6,6 +6,7 @@ logs/ pipelines.json followups.json siphon-watermarks.json +identity-state.json review/ __pycache__/ diff --git a/bin/identity-broker.py b/bin/identity-broker.py new file mode 100755 index 0000000..0c7a648 --- /dev/null +++ b/bin/identity-broker.py @@ -0,0 +1,342 @@ +#!/usr/bin/env python3 +"""identity-broker.py — Scope lifecycle over proxy providers. + +Owns identity-state.json (gitignored runtime state: scope -> label +assignments, no secrets) and drives providers through it: + + up resolve fingerprint, provision its scope + down teardown scope network, drop assignment + cycle rotate to a fresh identity (bumps cycles) + exec -- run a command inside the scope network + routes read-only route/tunnel status + status scopes with emails/labels (no key material) + bind --from --session --fp + attribute live runs (agent-manager --once + --json) to a scope after verifying the + session exists in the scan + +v1 boundary: the broker scopes NETWORK identity only. It never sees, +stores, or prints key bytes — fingerprints and emails are the only +identifiers here. Short-lived brokered credential issuance is a +deferred stage; harnesses receive keys through existing means. +""" + +from __future__ import annotations + +import argparse +import importlib.util +import json +import sys +import time +from pathlib import Path +from typing import Any, Callable, Dict, List, Optional, Tuple + +REPO_ROOT = Path(__file__).resolve().parent.parent +BIN_DIR = REPO_ROOT / "bin" +STATE_FILE = REPO_ROOT / "identity-state.json" + +RunFn = Callable[..., Tuple[int, str]] + + +def _load(name: str, modname: str): + path = BIN_DIR / name + spec = importlib.util.spec_from_file_location(modname, path) + mod = importlib.util.module_from_spec(spec) + sys.modules[modname] = mod + spec.loader.exec_module(mod) + return mod + + +resolve_mod = _load("identity-resolve.py", "identity_resolve") +provider_mod = _load("identity-provider.py", "identity_provider") + +PROVIDERS = {"warp": provider_mod.WarpProvider(), + "wireguard": provider_mod.GenericWireGuardProvider(), + "socks": provider_mod.SocksProxyProvider()} + + +class BrokerError(RuntimeError): + """A broker operation failed (message is safe to show).""" + + +def load_state(path: str | Path = STATE_FILE) -> Dict[str, Any]: + """Load broker state. Missing/corrupt -> empty (never raises).""" + try: + with open(path, "r") as f: + data = json.load(f) + if isinstance(data, dict): + data.setdefault("scopes", {}) + data.setdefault("bindings", []) + return data + except Exception: + pass + return {"scopes": {}, "bindings": []} + + +def save_state(state: Dict[str, Any], + path: str | Path = STATE_FILE) -> None: + with open(path, "w") as f: + json.dump(state, f, indent=2, sort_keys=True) + + +def _provider(name: str = "warp"): + try: + prov = PROVIDERS[name] + except KeyError: + raise BrokerError("unknown provider %r (have: %s)" + % (name, ", ".join(sorted(PROVIDERS)))) + if not prov.ready: + raise BrokerError("provider %r is boilerplate (not implemented); " + "warp is the live provider" % (name,)) + return prov + + +def _scope_for_fp(fp: str, map_path=None) -> Dict[str, Any]: + scope = resolve_mod.resolve_scope( + resolve_mod.load_map(map_path or resolve_mod.MAP_FILE), fp) + if scope is None: + raise BrokerError("unknown fingerprint (not in identity-map.json)") + scope["slug"] = resolve_mod.scope_slug(scope) + return scope + + +def _unit_key(fp_or_unit: str, state: Dict[str, Any], + map_path=None) -> str: + """Resolve CLI input (fp or scope unit) to a state scopes key.""" + scopes = state.get("scopes", {}) + if fp_or_unit in scopes: + return fp_or_unit + try: + scope = _scope_for_fp(fp_or_unit, map_path) + except BrokerError: + raise BrokerError("no scope for %r (unknown fingerprint, no " + "such unit)" % (fp_or_unit,)) + if scope["unit"] not in scopes: + raise BrokerError("scope %s is not up" % (scope["unit"],)) + return scope["unit"] + + +def op_up(fp: str, run: Optional[RunFn] = None, map_path=None, + state_path: str | Path = STATE_FILE, + provider_name: str = "warp") -> Dict[str, Any]: + """Provision a scope's network. Idempotent (re-up returns existing).""" + scope = _scope_for_fp(fp, map_path) + state = load_state(state_path) + if scope["unit"] in state["scopes"]: + return {"ok": "exists", **state["scopes"][scope["unit"]]} + label = scope["slug"] + try: + res = _provider(provider_name).provision(label, run=run) + except provider_mod.ProviderError as e: + raise BrokerError(str(e)) + rec = {"scope": scope["scope"], "email": scope["email"], + "origins": scope["origins"], "label": label, + "provider": provider_name, "netns": res.get("netns", ""), + "created": int(time.time()), "cycles": 0} + state["scopes"][scope["unit"]] = rec + save_state(state, state_path) + return {"ok": "true", **rec} + + +def op_down(fp_or_unit: str, run: Optional[RunFn] = None, map_path=None, + state_path: str | Path = STATE_FILE) -> Dict[str, Any]: + """Teardown a scope's network and drop its assignment + bindings.""" + state = load_state(state_path) + unit = _unit_key(fp_or_unit, state, map_path) + rec = state["scopes"][unit] + try: + _provider(rec.get("provider", "warp")).teardown(rec["label"], + run=run) + except provider_mod.ProviderError as e: + raise BrokerError(str(e)) + del state["scopes"][unit] + state["bindings"] = [b for b in state.get("bindings", []) + if b.get("unit") != unit] + save_state(state, state_path) + return {"ok": "true", "unit": unit, "label": rec["label"]} + + +def op_cycle(fp: str, run: Optional[RunFn] = None, map_path=None, + state_path: str | Path = STATE_FILE) -> Dict[str, Any]: + """Rotate a scope to a fresh identity (bumps the cycle count).""" + scope = _scope_for_fp(fp, map_path) + state = load_state(state_path) + if scope["unit"] not in state["scopes"]: + raise BrokerError("scope %s is not up (up it first)" + % (scope["unit"],)) + rec = state["scopes"][scope["unit"]] + try: + _provider(rec.get("provider", "warp")).cycle(rec["label"], + run=run) + except provider_mod.ProviderError as e: + raise BrokerError(str(e)) + rec["cycles"] = int(rec.get("cycles", 0)) + 1 + save_state(state, state_path) + return {"ok": "true", "unit": scope["unit"], "label": rec["label"], + "cycles": rec["cycles"]} + + +def op_exec(fp: str, cmd: List[str], run: Optional[RunFn] = None, + map_path=None, + state_path: str | Path = STATE_FILE) -> Tuple[int, str]: + """Run cmd inside the scope's network. Returns (rc, output).""" + scope = _scope_for_fp(fp, map_path) + state = load_state(state_path) + if scope["unit"] not in state["scopes"]: + raise BrokerError("scope %s is not up (up it first)" + % (scope["unit"],)) + rec = state["scopes"][scope["unit"]] + try: + return _provider(rec.get("provider", "warp")).exec( + rec["label"], cmd, run=run) + except provider_mod.ProviderError as e: + raise BrokerError(str(e)) + + +def op_routes(fp: str, run: Optional[RunFn] = None, map_path=None, + state_path: str | Path = STATE_FILE) -> Dict[str, Any]: + """Read-only route/tunnel status for a scope.""" + scope = _scope_for_fp(fp, map_path) + state = load_state(state_path) + if scope["unit"] not in state["scopes"]: + raise BrokerError("scope %s is not up (up it first)" + % (scope["unit"],)) + rec = state["scopes"][scope["unit"]] + try: + info = _provider(rec.get("provider", "warp")).routes( + rec["label"], run=run) + except provider_mod.ProviderError as e: + raise BrokerError(str(e)) + return {"unit": scope["unit"], "email": scope["email"], **info} + + +def op_status(run: Optional[RunFn] = None, + state_path: str | Path = STATE_FILE) -> Dict[str, Any]: + """Scopes with live provider status. Emails/labels only.""" + state = load_state(state_path) + scopes = [] + for unit, rec in sorted(state.get("scopes", {}).items()): + try: + live = _provider(rec.get("provider", "warp")).status( + rec["label"], run=run) + except provider_mod.ProviderError as e: + live = {"conf": "?", "netns": "?", "egress": "n/a", + "error": str(e)} + scopes.append({"unit": unit, "email": rec.get("email", ""), + "scope": rec.get("scope", ""), + "label": rec.get("label", ""), + "provider": rec.get("provider", ""), + "cycles": rec.get("cycles", 0), **live}) + return {"scopes": scopes, "bindings": state.get("bindings", [])} + + +def op_bind(runs_path: str, session: str, fp: str, map_path=None, + state_path: str | Path = STATE_FILE) -> Dict[str, Any]: + """Attribute live runs to a scope, verifying against a scan. + + runs_path is agent-manager.py --once --json output. Every run + whose session group contains `session` is bound to the fp's scope + (fp must resolve; the scope need not be up — binding is + attribution, not network). Sessions absent from the scan are + refused (never bind what we cannot observe). + """ + scope = _scope_for_fp(fp, map_path) + try: + with open(runs_path, "r") as f: + scan = json.load(f) + runs = scan.get("runs", []) + if not isinstance(runs, list): + raise ValueError("no runs list") + except Exception as e: + raise BrokerError("cannot read runs scan %s: %s" % (runs_path, e)) + matched = [] + for r in runs: + if not isinstance(r, dict): + continue + group = str(r.get("session", "")).split(",") + if session in group: + matched.append(r) + if not matched: + raise BrokerError("session %r not observed in %s (refusing to " + "bind unseen runs)" % (session, runs_path)) + state = load_state(state_path) + now = int(time.time()) + new = [] + for r in matched: + rec = {"unit": scope["unit"], "email": scope["email"], + "device": r.get("device", ""), "type": r.get("type", ""), + "session": session, "pane": r.get("pane", ""), + "pid": r.get("pid", 0), "bound_at": now} + new.append(rec) + # Replace prior bindings for these exact runs (re-bind refreshes). + keys = {(b["device"], b.get("pane"), b.get("pid")) for b in new} + state["bindings"] = [b for b in state.get("bindings", []) + if (b.get("device"), b.get("pane"), + b.get("pid")) not in keys] + new + save_state(state, state_path) + return {"ok": "true", "unit": scope["unit"], "bound": len(new), + "runs": new} + + +def main(argv: Optional[List[str]] = None) -> int: + ap = argparse.ArgumentParser(prog="identity-broker.py") + ap.add_argument("--map", default=str(resolve_mod.MAP_FILE), + help="identity map (default: identity-map.json)") + ap.add_argument("--state", default=str(STATE_FILE), + help="broker state file (default: identity-state.json)") + sub = ap.add_subparsers(dest="cmd", required=True) + p = sub.add_parser("up", help="provision a scope network") + p.add_argument("fp") + p.add_argument("--provider", default="warp") + p = sub.add_parser("down", help="teardown a scope network") + p.add_argument("fp_or_unit") + p = sub.add_parser("cycle", help="rotate a scope identity") + p.add_argument("fp") + p = sub.add_parser("exec", help="run a command in a scope network") + p.add_argument("fp") + p.add_argument("exec_cmd", nargs=argparse.REMAINDER, + help="command (after --)") + p = sub.add_parser("routes", help="route/tunnel status for a scope") + p.add_argument("fp") + sub.add_parser("status", help="scopes + bindings (emails only)") + p = sub.add_parser("bind", help="attribute live runs to a scope") + p.add_argument("--from", dest="runs", required=True) + p.add_argument("--session", required=True) + p.add_argument("--fp", required=True) + args = ap.parse_args(argv) + mp, sp = args.map, args.state + try: + if args.cmd == "up": + print(json.dumps(op_up(args.fp, map_path=mp, + state_path=sp, + provider_name=args.provider), + indent=2)) + elif args.cmd == "down": + print(json.dumps(op_down(args.fp_or_unit, map_path=mp, + state_path=sp), indent=2)) + elif args.cmd == "cycle": + print(json.dumps(op_cycle(args.fp, map_path=mp, + state_path=sp), indent=2)) + elif args.cmd == "exec": + cmd = [c for c in (args.exec_cmd or []) if c != "--"] + rc, out = op_exec(args.fp, cmd, map_path=mp, + state_path=sp) + sys.stdout.write(out + ("\n" if out else "")) + return rc + elif args.cmd == "routes": + print(json.dumps(op_routes(args.fp, map_path=mp, + state_path=sp), indent=2)) + elif args.cmd == "status": + print(json.dumps(op_status(state_path=sp), indent=2)) + elif args.cmd == "bind": + print(json.dumps(op_bind(args.runs, args.session, args.fp, + map_path=mp, state_path=sp), + indent=2)) + return 0 + except BrokerError as e: + print("error: %s" % e) + return 1 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/bin/identity-provider.py b/bin/identity-provider.py new file mode 100755 index 0000000..da9c6d3 --- /dev/null +++ b/bin/identity-provider.py @@ -0,0 +1,237 @@ +#!/usr/bin/env python3 +"""identity-provider.py — Proxy provider implementations. + +A provider owns one network-identity substrate behind a fixed +interface: provision / teardown / cycle / exec / routes / status. +All subprocesses go through an injectable run function (same seam as +box-fleet-tui gather_*), so command shapes are unit-testable and no +test touches netns, sudo, or /etc/netvm. + +Security boundaries (from the repo's own scripts): +- Warp identities generate via netvm-new-identity.sh, which the user + explicitly authorized operators to run (see netvm-provision-node.sh + header). Generation installs a root-0600 conf and prints nothing. +- This code NEVER reads /etc/netvm and never prints key material. + Confs are consumed only by root tools (wg setconf inside netns). +- CLI-facing output carries emails, labels, and fingerprints only. +""" + +from __future__ import annotations + +import os +import re +import subprocess +import sys +from pathlib import Path +from typing import Callable, Dict, List, Optional, Tuple + +REPO_ROOT = Path(__file__).resolve().parent.parent +BIN_DIR = REPO_ROOT / "bin" + +RunFn = Callable[..., Tuple[int, str]] + +LABEL_RE = re.compile(r"^[a-z0-9][a-z0-9-]{0,22}$") + + +def _run(cmd: List[str], timeout: int = 120) -> Tuple[int, str]: + """Run cmd, capture output. Returns (returncode, combined_output).""" + try: + r = subprocess.run(cmd, capture_output=True, text=True, + timeout=timeout) + return r.returncode, ((r.stdout or "") + (r.stderr or "")).strip() + except subprocess.TimeoutExpired: + return 124, "timed out after %ds: %s" % (timeout, " ".join(cmd)) + except OSError as e: + return 127, str(e) + + +class ProviderError(RuntimeError): + """A provider operation failed (message is safe to show).""" + + +def check_label(label: str) -> str: + """Validate a netvm label. Returns it or raises ProviderError.""" + if not LABEL_RE.match(label or ""): + raise ProviderError( + "invalid label %r: lowercase letters, digits, hyphens " + "(max 23 chars)" % (label,)) + return label + + +class Provider: + """Interface every proxy provider implements. Boilerplate subclasses + override these with real substrate calls; see WarpProvider.""" + + name = "base" + ready = False + + def provision(self, label: str, + run: Optional[RunFn] = None) -> Dict[str, str]: + """Create the network identity + bring it up. Idempotent.""" + raise NotImplementedError + + def teardown(self, label: str, + run: Optional[RunFn] = None) -> Dict[str, str]: + """Bring the identity's network down (keeps the identity).""" + raise NotImplementedError + + def cycle(self, label: str, + run: Optional[RunFn] = None) -> Dict[str, str]: + """Rotate to a fresh identity (teardown + new identity + up).""" + raise NotImplementedError + + def exec(self, label: str, cmd: List[str], + run: Optional[RunFn] = None) -> Tuple[int, str]: + """Run cmd inside the identity's network. Returns (rc, output).""" + raise NotImplementedError + + def routes(self, label: str, + run: Optional[RunFn] = None) -> Dict[str, str]: + """Read-only route/tunnel status for the identity.""" + raise NotImplementedError + + def status(self, label: str, + run: Optional[RunFn] = None) -> Dict[str, str]: + """Read-only liveness: conf present, netns up, egress IP.""" + raise NotImplementedError + + +class WarpProvider(Provider): + """Cloudflare Warp provider on the established warp-* structures. + + Identity: /etc/netvm/