mirror of
https://scm.tikali.ai/tikali/applications/monky/monky-deployd.git
synced 2026-09-18 07:16:16 +00:00
31586c3058
install.sh wrote tenancy.port 8081 (the in-pod port) but the monky.tenancy.deploy intercept is 443, so the SDK found no intercept and the check-in failed with 'service not available' followed by the SDK's bare TypeError. Also wrap non-OSError SDK failures in a TransportError that names the host:port and what to check. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01KLB7jieMNRkTsJ2epr4Ds1
230 lines
8.3 KiB
Python
230 lines
8.3 KiB
Python
"""How the agent reaches the mesh, and a tiny HTTP client on top of it.
|
|
|
|
Three transports, one interface (`connect(host, port) -> socket`):
|
|
|
|
* `sdk` — the OpenZiti Python SDK dials the ziti service by its intercept name with the
|
|
box's own host identity (no tun, no root). Coexists with `ziti-edge-tunnel run-host`.
|
|
* `proxy` — `monky-deployd-proxy.service` runs `ziti tunnel proxy … monky.tenancy.deploy:18443
|
|
openbao:18200` as user ziti; the agent talks to 127.0.0.1:<port>. TLS SNI and
|
|
certificate checks still use the real hostname.
|
|
* `system` — plain DNS/TCP, for a laptop whose tunneler runs in `run` mode (tun + DNS).
|
|
|
|
Every network failure surfaces as `TransportError` (exit 75: temporary, retry next tick)."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import http.client
|
|
import json
|
|
import logging
|
|
import socket
|
|
import ssl
|
|
from dataclasses import dataclass
|
|
|
|
from monky_deployd import __version__
|
|
from monky_deployd.config import Config
|
|
|
|
log = logging.getLogger("monky-deployd.transport")
|
|
|
|
USER_AGENT = f"monky-deployd/{__version__}"
|
|
|
|
|
|
class TransportError(Exception):
|
|
"""A network-level failure: DNS, connect, TLS, timeout, reset. Retry next tick."""
|
|
|
|
|
|
class Transport:
|
|
name = "base"
|
|
|
|
def connect(self, host: str, port: int, timeout: float) -> socket.socket: # pragma: no cover
|
|
raise NotImplementedError
|
|
|
|
def describe(self) -> str:
|
|
return self.name
|
|
|
|
|
|
class SystemTransport(Transport):
|
|
name = "system"
|
|
|
|
def connect(self, host: str, port: int, timeout: float) -> socket.socket:
|
|
return socket.create_connection((host, port), timeout=timeout)
|
|
|
|
|
|
class ProxyTransport(Transport):
|
|
"""(host, port) -> 127.0.0.1:<proxy port>; anything unmapped is refused (no leaks)."""
|
|
|
|
name = "proxy"
|
|
|
|
def __init__(self, mapping: dict[tuple[str, int], tuple[str, int]]):
|
|
self.mapping = mapping
|
|
|
|
def connect(self, host: str, port: int, timeout: float) -> socket.socket:
|
|
try:
|
|
target = self.mapping[(host, port)]
|
|
except KeyError as exc:
|
|
raise TransportError(f"no proxy mapping for {host}:{port}") from exc
|
|
return socket.create_connection(target, timeout=timeout)
|
|
|
|
def describe(self) -> str:
|
|
return "proxy(" + ", ".join(f"{h}:{p}->{t[0]}:{t[1]}" for (h, p), t in self.mapping.items()) + ")"
|
|
|
|
|
|
class SdkTransport(Transport):
|
|
"""`import openziti` is deferred so the other transports work without the wheel."""
|
|
|
|
name = "sdk"
|
|
|
|
def __init__(self, identity_path: str):
|
|
self.identity_path = identity_path
|
|
self._ctx = None
|
|
|
|
def _load(self):
|
|
if self._ctx is not None:
|
|
return
|
|
try:
|
|
import openziti # type: ignore
|
|
except ImportError as exc: # pragma: no cover - exercised via a fake module in tests
|
|
raise TransportError(
|
|
"transport sdk: the openziti module is not installed in this venv; "
|
|
"use transport: proxy (monky-deployd-proxy.service) or system"
|
|
) from exc
|
|
try:
|
|
self._ctx = openziti.load(self.identity_path)
|
|
except Exception as exc:
|
|
raise TransportError(f"transport sdk: cannot load identity {self.identity_path}: {exc}") from exc
|
|
self._openziti = openziti
|
|
|
|
def connect(self, host: str, port: int, timeout: float) -> socket.socket:
|
|
self._load()
|
|
# monkeypatch() swaps socket.socket for the SDK's ZitiSocket for the duration of the
|
|
# block: an address that matches a ziti intercept is dialled over the mesh, anything
|
|
# else falls through to a plain socket (which the proxy transport would have refused —
|
|
# for sdk that is what we want: bao.cbs.tikali.net is an intercept, not public DNS).
|
|
try:
|
|
with self._openziti.monkeypatch():
|
|
return socket.create_connection((host, port), timeout=timeout)
|
|
except OSError as exc:
|
|
raise TransportError(f"transport sdk: dial {host}:{port} failed: {exc}") from exc
|
|
except Exception as exc: # noqa: BLE001 - the SDK raises bare Exception((code, msg)) and TypeError
|
|
# openziti-sdk-py: an address with NO matching intercept falls through to
|
|
# PySocket.connect(tuple) → TypeError; a matching intercept the identity may not dial
|
|
# raises Exception((-18, 'service not available')) — env-qa-02 pilot, 2026-09-07.
|
|
raise TransportError(
|
|
f"transport sdk: dial {host}:{port} failed: {exc} — no intercept for that host:port, or this "
|
|
"identity has no dial policy for the service (check the intercept port and the identity's attrs)"
|
|
) from exc
|
|
|
|
def describe(self) -> str:
|
|
return f"sdk(identity={self.identity_path})"
|
|
|
|
|
|
def build(cfg: Config) -> Transport:
|
|
if cfg.transport == "system":
|
|
return SystemTransport()
|
|
_, bao_host, bao_port = cfg.bao_url
|
|
if cfg.transport == "proxy":
|
|
|
|
def _addr(s: str) -> tuple[str, int]:
|
|
h, _, p = s.rpartition(":")
|
|
return h.strip("[]") or "127.0.0.1", int(p)
|
|
|
|
return ProxyTransport(
|
|
{
|
|
(cfg.tenancy.host, cfg.tenancy.port): _addr(cfg.tenancy.proxy_addr),
|
|
(bao_host, bao_port): _addr(cfg.bao.proxy_addr),
|
|
}
|
|
)
|
|
return SdkTransport(cfg.identity)
|
|
|
|
|
|
# --- HTTP ---------------------------------------------------------------------------------------
|
|
|
|
|
|
@dataclass
|
|
class HttpResponse:
|
|
status: int
|
|
headers: dict[str, str]
|
|
body: bytes
|
|
|
|
def json(self):
|
|
if not self.body:
|
|
return None
|
|
try:
|
|
return json.loads(self.body.decode())
|
|
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
|
|
raise TransportError(f"non-JSON response ({self.status})") from exc
|
|
|
|
def text(self) -> str:
|
|
return self.body.decode(errors="replace")
|
|
|
|
|
|
class _Conn(http.client.HTTPConnection):
|
|
"""http.client with the TCP connect delegated to a Transport (+ optional TLS with SNI)."""
|
|
|
|
def __init__(self, transport: Transport, host: str, port: int, timeout: float, ctx: ssl.SSLContext | None):
|
|
super().__init__(host, port, timeout=timeout)
|
|
self._transport = transport
|
|
self._ctx = ctx
|
|
|
|
def connect(self) -> None:
|
|
sock = self._transport.connect(self.host, self.port, self.timeout)
|
|
if self._ctx is not None:
|
|
sock = self._ctx.wrap_socket(sock, server_hostname=self.host)
|
|
self.sock = sock
|
|
|
|
|
|
class HttpClient:
|
|
def __init__(
|
|
self,
|
|
transport: Transport,
|
|
scheme: str,
|
|
host: str,
|
|
port: int,
|
|
*,
|
|
ca_bundle: str | None = None,
|
|
timeout: float = 30,
|
|
):
|
|
self.transport = transport
|
|
self.scheme = scheme
|
|
self.host = host
|
|
self.port = port
|
|
self.timeout = timeout
|
|
self.ctx: ssl.SSLContext | None = None
|
|
if scheme == "https":
|
|
self.ctx = ssl.create_default_context(cafile=ca_bundle) if ca_bundle else ssl.create_default_context()
|
|
self.ctx.minimum_version = ssl.TLSVersion.TLSv1_2
|
|
|
|
@property
|
|
def base(self) -> str:
|
|
return f"{self.scheme}://{self.host}:{self.port}"
|
|
|
|
def request(
|
|
self,
|
|
method: str,
|
|
path: str,
|
|
*,
|
|
json_body=None,
|
|
headers: dict[str, str] | None = None,
|
|
body: bytes | None = None,
|
|
) -> HttpResponse:
|
|
hdrs = {"User-Agent": USER_AGENT, "Accept": "application/json, application/x-tar;q=0.9, */*;q=0.1"}
|
|
if headers:
|
|
hdrs.update(headers)
|
|
if json_body is not None:
|
|
body = json.dumps(json_body, separators=(",", ":")).encode()
|
|
hdrs["Content-Type"] = "application/json"
|
|
conn = _Conn(self.transport, self.host, self.port, self.timeout, self.ctx)
|
|
try:
|
|
conn.request(method, path, body=body, headers=hdrs)
|
|
resp = conn.getresponse()
|
|
data = resp.read()
|
|
return HttpResponse(resp.status, {k.lower(): v for k, v in resp.getheaders()}, data)
|
|
except TransportError:
|
|
raise
|
|
except (OSError, http.client.HTTPException, ssl.SSLError) as exc:
|
|
raise TransportError(f"{method} {self.base}{path}: {exc.__class__.__name__}: {exc}") from exc
|
|
finally:
|
|
try:
|
|
conn.close()
|
|
except Exception: # pragma: no cover
|
|
pass
|