41499e069d
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.
343 lines
14 KiB
Python
Executable File
343 lines
14 KiB
Python
Executable File
#!/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 <fp> resolve fingerprint, provision its scope
|
|
down <fp|unit> teardown scope network, drop assignment
|
|
cycle <fp> rotate to a fresh identity (bumps cycles)
|
|
exec <fp> -- <cmd> run a command inside the scope network
|
|
routes <fp> read-only route/tunnel status
|
|
status scopes with emails/labels (no key material)
|
|
bind --from <runs.json> --session <s> --fp <f>
|
|
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())
|