from __future__ import annotations
import json
import logging
import os
import threading
from typing import Optional, Tuple
logger = logging.getLogger(__name__)
_DEFAULT_PSK = "00" * 32
_lock = threading.Lock()
_state: Optional[Tuple[object, object, str, object]] = None
_operator: Optional[Tuple[object, object]] = None
_renewal: Optional[object] = None
_GRANT_TTL_SECONDS = 365 * 24 * 60 * 60
def _config() -> dict:
return {
"bind_addr": os.environ.get("NET_MESH_BIND", "127.0.0.1:0"),
"psk": os.environ.get("NET_MESH_PSK", _DEFAULT_PSK),
"peers": os.environ.get("NET_MESH_PEERS", "").strip(),
"pin_store": (os.environ.get("NET_MESH_PIN_STORE") or "").strip() or None,
"identity_seed": (os.environ.get("NET_MESH_IDENTITY_SEED") or "").strip() or None,
}
def _build() -> Tuple[object, object, str, object]:
from net import NetMesh
from net_sdk import AsyncCapabilityGateway, default_pin_store_path
cfg = _config()
seed = None
if cfg["identity_seed"]:
try:
seed = bytes.fromhex(cfg["identity_seed"])
except ValueError as e:
raise RuntimeError(
"net plugin: NET_MESH_IDENTITY_SEED must be a 32-byte identity "
f"seed as 64 hex chars; got an unparseable value ({e})"
) from e
delegation = None
if seed is not None:
try:
from .delegation import GatewayDelegation
delegation = GatewayDelegation(seed)
except ImportError as e:
logger.info(
"net plugin: delegation surface unavailable (%s); running un-delegated", e
)
mesh = NetMesh(cfg["bind_addr"], cfg["psk"], identity_seed=seed, permissive_channels=True)
if cfg["peers"]:
try:
peers = json.loads(cfg["peers"])
except json.JSONDecodeError as e:
logger.warning("net plugin: NET_MESH_PEERS is not valid JSON (%s); no peers", e)
peers = []
for p in peers:
try:
mesh.connect(str(p["addr"]), str(p["pubkey"]), int(p["node_id"]))
except Exception as e: addr = p.get("addr") if isinstance(p, dict) else p
logger.warning("net plugin: connect to peer %s failed: %s", addr, e)
mesh.start()
try:
pin_store = cfg["pin_store"] or default_pin_store_path()
if not pin_store:
raise RuntimeError(
"net plugin: no pin-store path could be resolved; set NET_MESH_PIN_STORE"
)
if delegation is not None:
gateway = AsyncCapabilityGateway(
mesh,
pin_store_path=pin_store,
delegation_leaf=delegation.gateway_identity,
delegation_chain=delegation.chain_bytes(),
)
logger.info(
"net plugin: gateway delegation acquired (machine=%s, gateway=0x%s)",
delegation._machine_label,
delegation.gateway_id.hex()[:16],
)
else:
gateway = AsyncCapabilityGateway(mesh, pin_store_path=pin_store)
except BaseException:
try:
mesh.shutdown()
except Exception: logger.debug(
"net plugin: mesh shutdown during init rollback failed", exc_info=True
)
raise
try:
_start_device_renewal(mesh)
except Exception: logger.warning("net plugin: device-renewal setup failed", exc_info=True)
logger.info(
"net plugin: node up (id=%s, bind=%s, pin_store=%s, delegated=%s)",
getattr(mesh, "node_id", "?"),
cfg["bind_addr"],
pin_store,
delegation is not None,
)
return (mesh, gateway, pin_store, delegation)
def get_state() -> Tuple[object, object, str, object]:
global _state
if _state is not None:
return _state
with _lock:
if _state is None: _state = _build()
return _state
def check_net_available() -> bool:
try:
state = get_state()
except Exception as e: logger.debug("net plugin unavailable: %s", e)
return False
delegation = state[3]
if delegation is not None:
try:
if not delegation.verify():
logger.info(
"net plugin: gateway delegation invalid (revoked/expired); "
"tools unavailable"
)
return False
except Exception as e: logger.debug("net plugin: delegation verify error: %s", e)
return False
return True
def gateway():
return get_state()[1]
def pin_store_path() -> str:
return get_state()[2]
def delegation():
return get_state()[3]
def delegation_valid_for_invoke() -> bool:
try:
d = get_state()[3]
except Exception: return True
if d is None:
return True
try:
return d.verify()
except Exception: return False
def _env_int(name: str, default: int) -> int:
raw = (os.environ.get(name) or "").strip()
if not raw:
return default
try:
return int(raw)
except ValueError:
logger.warning("net plugin: %s is not an integer (%r); using default", name, raw)
return default
def _probe_writable(path: str) -> None:
parent = os.path.dirname(path) or "."
os.makedirs(parent, exist_ok=True)
probe = path + ".probe"
try:
with open(probe, "wb") as f:
f.write(b"net-mesh enrollment write probe")
finally:
try:
os.remove(probe)
except OSError:
pass
def _start_device_renewal(mesh_handle) -> None:
global _renewal
path = (os.environ.get("NET_MESH_DEVICE_ENROLLMENT") or "").strip() or None
if not path:
return
import time
from net import DeviceEnrollment, Identity, InviteToken
from .renewal import RenewalService
enrollment = DeviceEnrollment.load(path) if enrollment is None:
invite = (os.environ.get("NET_MESH_INVITE") or "").strip() or None
if not invite:
logger.info(
"net plugin: no device enrollment at %s and no NET_MESH_INVITE; "
"device renewal idle",
path,
)
return
_probe_writable(path)
device = Identity.generate()
name = (os.environ.get("NET_MESH_DEVICE_NAME") or "").strip() or "device"
chain = mesh_handle.join(device, invite, name, [])
rendezvous = InviteToken.decode(invite).rendezvous
enrollment = DeviceEnrollment(device, chain, rendezvous, int(time.time()))
save_error = None
for _ in range(3):
try:
enrollment.save(path)
save_error = None
break
except Exception as e: save_error = e
time.sleep(0.2)
if save_error is None:
logger.info("net plugin: enrolled as a new device (persisted to %s)", path)
else:
logger.error(
"net plugin: enrolled as a new device but persisting to %s "
"failed (%s); the enrollment is only in memory — a restart "
"before a successful save will need a fresh invite",
path,
save_error,
)
svc = RenewalService(
mesh_handle,
enrollment,
path,
check_interval=_env_int("NET_MESH_RENEWAL_INTERVAL", 24 * 60 * 60),
renewal_window=_env_int("NET_MESH_RENEWAL_WINDOW", 30 * 24 * 60 * 60),
)
svc.start()
_renewal = svc
logger.info("net plugin: silent device renewal started")
def mesh():
return get_state()[0]
def _build_operator(mesh_handle) -> Tuple[object, object]:
from net import Identity, OperatorEnrollment
cfg = _config()
if not cfg["identity_seed"]:
raise RuntimeError(
"net plugin: mesh device enrollment needs the user root identity; "
"set NET_MESH_IDENTITY_SEED"
)
root = Identity.from_seed(bytes.fromhex(cfg["identity_seed"]))
dev = (os.environ.get("NET_MESH_DEVICE_STORE") or "").strip() or None
rev = (os.environ.get("NET_MESH_REVOCATION_STORE") or "").strip() or None
if dev and rev:
operator = OperatorEnrollment(root, dev, rev)
elif dev or rev:
missing = "NET_MESH_REVOCATION_STORE" if dev else "NET_MESH_DEVICE_STORE"
present = "NET_MESH_DEVICE_STORE" if dev else "NET_MESH_REVOCATION_STORE"
raise RuntimeError(
f"net plugin: {present} is set but {missing} is not — the store "
"overrides must be set together (mixing an override with the "
"machine-shared default would split the inventory)"
)
else:
operator = OperatorEnrollment.with_default_paths(root)
handle = mesh_handle.serve_enrollment_auto(operator, _GRANT_TTL_SECONDS)
return (operator, handle)
def operator():
global _operator
if _operator is not None:
return _operator[0]
if not _config()["identity_seed"]:
raise RuntimeError(
"net plugin: mesh device enrollment needs the user root identity; "
"set NET_MESH_IDENTITY_SEED"
)
mesh_handle = mesh()
with _lock:
if _operator is None:
_operator = _build_operator(mesh_handle)
return _operator[0]
def shutdown() -> None:
global _state, _operator, _renewal
with _lock:
op_handle = _operator[1] if _operator is not None else None
_operator = None
renewal = _renewal
_renewal = None
mesh_handle = _state[0] if _state is not None else None
_state = None
if renewal is not None:
try:
renewal.stop()
except Exception: logger.debug("net plugin: renewal stop failed", exc_info=True)
if op_handle is not None:
try:
op_handle.stop()
except Exception: logger.debug("net plugin: enrollment serve-handle stop failed", exc_info=True)
if mesh_handle is not None:
try:
mesh_handle.shutdown()
except Exception as e: logger.debug("net plugin: shutdown error: %s", e)