refactor(relay): unify env/cfg URL readers and connector JSON POST; compact media/transport/manifest/__init__ docstrings
This commit is contained in:
+60
-88
@@ -1,15 +1,12 @@
|
||||
"""Relay/connector support package for the Hermes gateway.
|
||||
|
||||
EXPERIMENTAL. Gateway side of the "Gateway Gateway" relay design: a generic
|
||||
EXPERIMENTAL gateway side of the "Gateway Gateway" relay design: a generic
|
||||
``RelayAdapter`` plus the wire-serializable ``CapabilityDescriptor`` the connector
|
||||
hands it at handshake time, and the production ``WebSocketRelayTransport`` that
|
||||
dials the connector. The public API (module names, descriptor field set, transport
|
||||
protocol) MAY CHANGE without a deprecation cycle until at least two real Class-1
|
||||
platforms have shaken out the schema. See ``docs/relay-connector-contract.md``.
|
||||
|
||||
Activation is config-driven, not a feature flag: the relay platform is registered
|
||||
when a connector relay URL is configured (``GATEWAY_RELAY_URL`` env or
|
||||
``gateway.relay_url`` in config.yaml) — the same shape as ``gateway.proxy_url``.
|
||||
hands it at handshake, and the production ``WebSocketRelayTransport``. The public
|
||||
API MAY CHANGE without a deprecation cycle until >=2 real Class-1 platforms have
|
||||
shaken out the schema (``docs/relay-connector-contract.md``). Activation is
|
||||
config-driven: the relay platform is registered when a connector relay URL is set
|
||||
(``GATEWAY_RELAY_URL`` env or ``gateway.relay_url``), like ``gateway.proxy_url``.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -59,36 +56,32 @@ def _gateway_cfg() -> dict:
|
||||
|
||||
|
||||
def _env_or_cfg(env_var: str, cfg_key: str) -> str:
|
||||
"""Env var first (Docker/NAS stamp), then ``gateway.<cfg_key>`` in config.yaml.
|
||||
|
||||
Returns the stripped value, ``""`` when neither is set.
|
||||
"""
|
||||
"""Env var first (Docker/NAS stamp), then ``gateway.<cfg_key>``; stripped, ``""`` when unset."""
|
||||
value = os.environ.get(env_var, "").strip()
|
||||
return value or str(_gateway_cfg().get(cfg_key, "") or "").strip()
|
||||
|
||||
|
||||
def relay_url() -> Optional[str]:
|
||||
"""The connector relay endpoint URL, or None when relay is not configured.
|
||||
def _env_or_cfg_url(env_var: str, cfg_key: str) -> Optional[str]:
|
||||
"""``_env_or_cfg`` for URL-shaped values: trailing slash stripped, None when unset."""
|
||||
return _env_or_cfg(env_var, cfg_key).rstrip("/") or None
|
||||
|
||||
``GATEWAY_RELAY_URL`` env, then ``gateway.relay_url``. A non-empty value
|
||||
activates the relay platform; absence means a normal direct gateway.
|
||||
"""
|
||||
return _env_or_cfg("GATEWAY_RELAY_URL", "relay_url").rstrip("/") or None
|
||||
|
||||
def relay_url() -> Optional[str]:
|
||||
"""The connector relay endpoint URL, or None. A non-empty value activates the relay platform."""
|
||||
return _env_or_cfg_url("GATEWAY_RELAY_URL", "relay_url")
|
||||
|
||||
|
||||
def relay_platform_identities() -> list[tuple[str, str]]:
|
||||
"""The ordered (platform, bot_id) pairs this gateway fronts over the relay.
|
||||
|
||||
One gateway fronts a SET of platforms on one WS connection, from the env-stamped
|
||||
deploy config: ``GATEWAY_RELAY_PLATFORMS`` (comma-sep, e.g. ``discord,telegram``)
|
||||
and ``GATEWAY_RELAY_BOT_IDS`` (JSON map ``{"discord": {"botId": "..."}, ...}``).
|
||||
The FIRST pair is the default the handshake/descriptor falls back to. A platform
|
||||
listed but absent from the ids map resolves with an empty bot_id (the connector
|
||||
rejects an unprovisioned platform with a structured failure). Defaults to
|
||||
``[("relay", "")]`` — the generic single-plane fallback — when nothing is set.
|
||||
deploy config: ``GATEWAY_RELAY_PLATFORMS`` (comma-sep) and ``GATEWAY_RELAY_BOT_IDS``
|
||||
(JSON map ``{"discord": {"botId": "..."}, ...}``). The FIRST pair is the default
|
||||
the handshake/descriptor falls back to. A platform absent from the ids map gets an
|
||||
empty bot_id (the connector rejects an unprovisioned platform with a structured
|
||||
failure). Defaults to ``[("relay", "")]`` — the generic single-plane fallback.
|
||||
"""
|
||||
platforms_raw = os.environ.get("GATEWAY_RELAY_PLATFORMS", "").strip()
|
||||
platforms = [p.strip() for p in platforms_raw.split(",") if p.strip()]
|
||||
platforms = [p.strip() for p in os.environ.get("GATEWAY_RELAY_PLATFORMS", "").split(",") if p.strip()]
|
||||
if not platforms:
|
||||
return [("relay", "")]
|
||||
ids = _relay_bot_ids_map()
|
||||
@@ -129,12 +122,9 @@ def relay_platform_identity() -> tuple[str, str]:
|
||||
|
||||
|
||||
def relay_connection_auth() -> tuple[Optional[str], Optional[str]]:
|
||||
"""The (gateway_id, upgrade_secret) this gateway authenticates the WS upgrade with.
|
||||
|
||||
Both come from enrollment (``GATEWAY_RELAY_ID`` / ``GATEWAY_RELAY_SECRET`` env,
|
||||
then ``gateway.relay_id`` / ``gateway.relay_secret``). Either absent ->
|
||||
``(None, None)`` and the transport dials unauthenticated.
|
||||
"""
|
||||
"""The (gateway_id, upgrade_secret) from enrollment (``GATEWAY_RELAY_ID`` /
|
||||
``GATEWAY_RELAY_SECRET``, then ``gateway.relay_id`` / ``gateway.relay_secret``).
|
||||
Either absent -> ``(None, None)`` and the transport dials unauthenticated."""
|
||||
gateway_id = _env_or_cfg("GATEWAY_RELAY_ID", "relay_id")
|
||||
secret = _env_or_cfg("GATEWAY_RELAY_SECRET", "relay_secret")
|
||||
return (gateway_id or None, secret or None)
|
||||
@@ -145,11 +135,10 @@ def relay_endpoint() -> Optional[str]:
|
||||
|
||||
The connector delivers signed inbound POSTs here and stores it on the tenant's
|
||||
route rows; gateway-asserted but scoped to the verified tenant, so a dishonest
|
||||
gateway can only misdirect its OWN inbound. ``GATEWAY_RELAY_ENDPOINT`` env
|
||||
(self-hosted or NAS-stamped), then ``gateway.relay_endpoint``. Absent -> the
|
||||
gateway provisions outbound-only (no inbound routes written).
|
||||
gateway can only misdirect its OWN inbound. ``GATEWAY_RELAY_ENDPOINT`` env, then
|
||||
``gateway.relay_endpoint``. Absent -> outbound-only (no inbound routes written).
|
||||
"""
|
||||
return _env_or_cfg("GATEWAY_RELAY_ENDPOINT", "relay_endpoint").rstrip("/") or None
|
||||
return _env_or_cfg_url("GATEWAY_RELAY_ENDPOINT", "relay_endpoint")
|
||||
|
||||
|
||||
def relay_route_keys() -> list[str]:
|
||||
@@ -172,11 +161,9 @@ def relay_route_keys() -> list[str]:
|
||||
def relay_instance_id() -> Optional[str]:
|
||||
"""Stable per-instance id forwarded at provision (``GATEWAY_RELAY_INSTANCE_ID`` / ``gateway.relay_instance_id``).
|
||||
|
||||
Binds the connector's ``gatewayId -> instanceId`` so inbound can route
|
||||
per-instance rather than tenant-broadcast. NAS stamps its ``AgentInstance.id``
|
||||
for a managed agent; a self-hosted operator may set it. Gateway-asserted but
|
||||
tenant-scoped like ``relay_endpoint()``. Absent -> the connector stores null
|
||||
(back-compat: no per-instance binding yet).
|
||||
Binds the connector's ``gatewayId -> instanceId`` so inbound can route per-instance
|
||||
rather than tenant-broadcast (NAS stamps its ``AgentInstance.id``). Gateway-asserted
|
||||
but tenant-scoped like ``relay_endpoint()``. Absent -> the connector stores null.
|
||||
"""
|
||||
return _env_or_cfg("GATEWAY_RELAY_INSTANCE_ID", "relay_instance_id") or None
|
||||
|
||||
@@ -184,14 +171,12 @@ def relay_instance_id() -> Optional[str]:
|
||||
def relay_wake_url() -> Optional[str]:
|
||||
"""The gateway's WAKE URL forwarded at provision (``GATEWAY_RELAY_WAKE_URL`` / ``gateway.relay_wake_url``).
|
||||
|
||||
A poke target the connector GETs (payload-free) when a going-idle destination for
|
||||
A payload-free poke target the connector GETs when a going-idle destination for
|
||||
this instance receives its first buffered event, so a suspended gateway wakes,
|
||||
reconnects its relay WS and drains its backlog. NAS stamps it for a managed
|
||||
container; a self-hosted operator sets it (or ``hermes gateway enroll --wake-url``).
|
||||
Tenant-scoped like ``relay_instance_id()``. Absent -> the connector can't wake
|
||||
this instance (buffering still works; the gateway drains on next reconnect).
|
||||
reconnects its relay WS and drains its backlog. Absent -> the connector can't
|
||||
wake this instance (buffering still works; the gateway drains on next reconnect).
|
||||
"""
|
||||
return _env_or_cfg("GATEWAY_RELAY_WAKE_URL", "relay_wake_url").rstrip("/") or None
|
||||
return _env_or_cfg_url("GATEWAY_RELAY_WAKE_URL", "relay_wake_url")
|
||||
|
||||
|
||||
def relay_display_name() -> Optional[str]:
|
||||
@@ -199,9 +184,8 @@ def relay_display_name() -> Optional[str]:
|
||||
|
||||
Primary source of the connector's multi-agent reply-attribution prefix
|
||||
(``**<displayName>:** ``). ``GATEWAY_RELAY_DISPLAY_NAME`` env, then the skin's
|
||||
branded agent name (the CLI banner value), so a skin rename propagates on the
|
||||
next boot's re-provision. Absent -> the connector falls back to the instance's
|
||||
linked-owner identity, else skips the prefix.
|
||||
branded agent name, so a skin rename propagates on the next boot's re-provision.
|
||||
Absent -> the connector falls back to the instance's linked-owner identity.
|
||||
"""
|
||||
value = os.environ.get("GATEWAY_RELAY_DISPLAY_NAME", "").strip()
|
||||
if not value:
|
||||
@@ -224,26 +208,20 @@ def relay_display_name() -> Optional[str]:
|
||||
# ─────────────────────────── connector HTTP ───────────────────────────
|
||||
|
||||
|
||||
def _http_base(relay_dial_url: str) -> str:
|
||||
"""``ws(s)://…/relay`` dial URL -> ``http(s)://…`` connector base (shared with media)."""
|
||||
def _connector_url(relay_dial_url: str, path: str) -> str:
|
||||
"""``ws(s)://…/relay`` dial URL -> ``http(s)://…{path}`` connector route."""
|
||||
from gateway.relay.media import media_base_url
|
||||
|
||||
return media_base_url(relay_dial_url)
|
||||
return f"{media_base_url(relay_dial_url)}{path}"
|
||||
|
||||
|
||||
def _provision_url(relay_dial_url: str) -> str:
|
||||
"""Map the dial URL to the ``/relay/provision`` POST URL."""
|
||||
return f"{_http_base(relay_dial_url)}/relay/provision"
|
||||
return _connector_url(relay_dial_url, "/relay/provision")
|
||||
|
||||
|
||||
def _policy_url(relay_dial_url: str) -> str:
|
||||
"""Map the dial URL to the ``/relay/policy`` POST URL (relevance-policy channel)."""
|
||||
return f"{_http_base(relay_dial_url)}/relay/policy"
|
||||
|
||||
|
||||
def _json_post_request(url: str, token: str, body: dict) -> urllib.request.Request:
|
||||
"""A bearer-authenticated JSON POST request to the connector."""
|
||||
return urllib.request.Request(
|
||||
def _json_post(url: str, token: str, body: dict, timeout: float):
|
||||
"""Bearer-authenticated JSON POST; returns the open response (caller reads it)."""
|
||||
req = urllib.request.Request(
|
||||
url,
|
||||
data=json.dumps(body).encode("utf-8"),
|
||||
method="POST",
|
||||
@@ -253,24 +231,22 @@ def _json_post_request(url: str, token: str, body: dict) -> urllib.request.Reque
|
||||
"Accept": "application/json",
|
||||
},
|
||||
)
|
||||
return urllib.request.urlopen(req, timeout=timeout)
|
||||
|
||||
|
||||
def relay_relevance_policy(platform: Optional[str] = None) -> Optional[dict]:
|
||||
"""Project a fronted platform's RELEVANCE config into the connector's generic vocabulary.
|
||||
|
||||
The connector's relevance gate reasons over a platform-agnostic policy, keyed by
|
||||
``(tenant, platform, instanceId)``. Mapping from the agent's existing knobs:
|
||||
- ``requireAddress`` <- the platform's ``require_mention``
|
||||
- ``freeResponseScopes`` <- the platform's ``free_response_channels``
|
||||
- ``allowOtherBots`` <- ``{PLATFORM}_ALLOW_BOTS`` in {"mentions","all"}
|
||||
|
||||
Read from the platform's config block (``discord:``), falling back to the bridged
|
||||
top-level keys, then env. ``platform`` defaults to the PRIMARY fronted platform.
|
||||
Returns None when relay isn't configured or nothing is CONFIGURED to declare, so
|
||||
the connector's default (mention-gated) applies. The condition is "require_mention
|
||||
is unset", NOT "falsy": an EXPLICIT ``require_mention: false`` is a non-default
|
||||
choice that MUST be declared or the connector would mention-gate an agent
|
||||
configured to free-respond.
|
||||
The connector's relevance gate reasons over a platform-agnostic policy keyed by
|
||||
``(tenant, platform, instanceId)``: ``requireAddress`` <- ``require_mention``,
|
||||
``freeResponseScopes`` <- ``free_response_channels``, ``allowOtherBots`` <-
|
||||
``{PLATFORM}_ALLOW_BOTS`` in {"mentions","all"}. Read from the platform's config
|
||||
block (``discord:``), falling back to the bridged top-level keys, then env.
|
||||
``platform`` defaults to the PRIMARY fronted platform. Returns None when relay
|
||||
isn't configured or nothing is CONFIGURED to declare, so the connector's default
|
||||
(mention-gated) applies. The condition is "require_mention is unset", NOT "falsy":
|
||||
an EXPLICIT ``require_mention: false`` is a non-default choice that MUST be
|
||||
declared or the connector would mention-gate an agent configured to free-respond.
|
||||
"""
|
||||
if platform is None:
|
||||
platform, _bot_id = relay_platform_identity()
|
||||
@@ -310,9 +286,6 @@ def relay_relevance_policy(platform: Optional[str] = None) -> Optional[dict]:
|
||||
allow_bots_env = os.environ.get(f"{platform.upper()}_ALLOW_BOTS", "").lower().strip()
|
||||
allow_other_bots = allow_bots_env in {"mentions", "all"}
|
||||
|
||||
# Nothing CONFIGURED to declare => keep the connector's default. The test is
|
||||
# "require_mention is unset", NOT "falsy": an EXPLICIT `require_mention: false`
|
||||
# is a non-default choice that MUST be declared (see docstring).
|
||||
if require_mention is None and not free_response and not allow_other_bots:
|
||||
return None
|
||||
return {
|
||||
@@ -356,9 +329,8 @@ def _post_provision(
|
||||
for key, value in (("instanceId", instance_id), ("wakeUrl", wake_url), ("displayName", display_name)):
|
||||
if value:
|
||||
body[key] = value
|
||||
req = _json_post_request(provision_url, access_token, body)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=timeout) as resp:
|
||||
with _json_post(provision_url, access_token, body, timeout) as resp:
|
||||
payload = json.loads(resp.read().decode())
|
||||
except urllib.error.HTTPError as exc:
|
||||
detail = ""
|
||||
@@ -382,9 +354,9 @@ def _resolve_relay_identity_token() -> str:
|
||||
|
||||
Canonical resolver shared by runtime self-provision and ``hermes gateway enroll``.
|
||||
Modes, in precedence order:
|
||||
1. Generic OIDC client-credentials (self-hosted IdP, no Nous Portal): when
|
||||
``gateway.idp.token_url`` (``GATEWAY_RELAY_IDP_TOKEN_URL``) is set together
|
||||
with client id + secret, POST the OAuth2 ``client_credentials`` grant.
|
||||
1. Generic OIDC client-credentials (self-hosted IdP): ``gateway.idp.token_url``
|
||||
(``GATEWAY_RELAY_IDP_TOKEN_URL``) set together with client id + secret ->
|
||||
POST the OAuth2 ``client_credentials`` grant.
|
||||
1b. Ambient token endpoint: ``token_url`` with NEITHER client_id nor
|
||||
client_secret is a metadata-server-style endpoint (e.g. Domino's
|
||||
``$DOMINO_API_PROXY/access-token``): a plain GET whose body IS the token,
|
||||
@@ -579,9 +551,8 @@ def _post_policy(*, policy_url: str, token: str, policy: dict, timeout: float =
|
||||
the connector resolves ``{tenant, instanceId}`` from its stored secret record,
|
||||
never the body. Raises RuntimeError on transport failure.
|
||||
"""
|
||||
req = _json_post_request(policy_url, token, policy)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=timeout) as resp:
|
||||
with _json_post(policy_url, token, policy, timeout) as resp:
|
||||
return int(resp.status)
|
||||
except urllib.error.HTTPError as exc:
|
||||
return int(exc.code)
|
||||
@@ -620,13 +591,14 @@ def send_relay_policy() -> bool:
|
||||
logger.warning("relay policy declaration failed to build token (%s); connector keeps prior policy", exc)
|
||||
return False
|
||||
|
||||
policy_url = _connector_url(dial_url, "/relay/policy")
|
||||
any_declared = False
|
||||
for platform, _bot_id in relay_platform_identities():
|
||||
policy = relay_relevance_policy(platform)
|
||||
if policy is None:
|
||||
continue
|
||||
try:
|
||||
status = _post_policy(policy_url=_policy_url(dial_url), token=token, policy=policy)
|
||||
status = _post_policy(policy_url=policy_url, token=token, policy=policy)
|
||||
except Exception as exc: # noqa: BLE001 - boot must survive a policy-declare failure
|
||||
logger.warning(
|
||||
"relay policy declaration failed for platform=%s (%s); continuing", platform, exc
|
||||
|
||||
@@ -1,20 +1,13 @@
|
||||
"""Gateway-declared slash-command manifest for the relay lane (Phase 4).
|
||||
"""Gateway-declared Discord slash-command manifest for the relay lane (Phase 4).
|
||||
|
||||
Over the relay the CONNECTOR holds the Discord token, so the gateway DECLARES
|
||||
its command set on the ``hello`` frame (``command_manifest``) and the connector
|
||||
reconciles Discord's global registration against it (GET → diff → bulk PUT,
|
||||
idempotent, best-effort). This module is the single source of truth for that
|
||||
declaration and MIRRORS the native tree (`_register_slash_commands`,
|
||||
plugins/platforms/discord/adapter.py) — same names, same descriptions.
|
||||
Interactions come back over the passthrough plane and are normalized by
|
||||
RelayAdapter._discord_interaction_to_event into the same "/name args" COMMAND
|
||||
events the dispatcher already routes, so declaring a command here needs NO new
|
||||
handler.
|
||||
|
||||
Wire shape per entry: {name, description, options?}; options rows are Discord
|
||||
option objects passed through verbatim. Names must satisfy Discord's CHAT_INPUT
|
||||
rules ([a-z0-9_-]{1,32}); the connector drops invalid entries (fail-open per
|
||||
entry, never the whole manifest).
|
||||
The CONNECTOR holds the Discord token, so the gateway declares its command set on
|
||||
the ``hello`` frame and the connector reconciles Discord's registration (idempotent,
|
||||
best-effort). MIRRORS the native tree (plugins/platforms/discord/adapter.py
|
||||
``_register_slash_commands``) — same names, same descriptions; interactions return
|
||||
via the passthrough plane as ordinary "/name args" COMMAND events, so a new entry
|
||||
needs NO new handler. Wire shape per entry: {name, description, options?} with
|
||||
Discord option objects verbatim; names must match ``[a-z0-9_-]{1,32}`` (the
|
||||
connector drops invalid entries, never the whole manifest).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
+16
-45
@@ -1,19 +1,12 @@
|
||||
"""Relay media client — gateway↔connector media plane (Phase 2). EXPERIMENTAL.
|
||||
|
||||
The relay wire carries media BY REFERENCE: inbound ``media_urls`` name
|
||||
connector re-hosts (``{connector}/relay/media/{id}``) and an outbound
|
||||
``send_media`` op names a ``source_url`` the connector resolves back to bytes.
|
||||
|
||||
- ``download(url)`` → GET a re-hosted attachment to a local temp file (the
|
||||
agent's vision/file tools consume LOCAL paths, like every native adapter).
|
||||
- ``upload(path)`` → POST local bytes to ``/relay/media``; returns the
|
||||
``/relay/media/{id}`` reference for a subsequent ``send_media`` op, so a
|
||||
locally-generated artifact crosses without the gateway needing a public URL.
|
||||
|
||||
Both present the same per-gateway signed bearer the WS upgrade uses
|
||||
(``make_upgrade_token``). Uploads are per-gateway-owned on the connector;
|
||||
downloads accept this gateway's uploads and connector ingest re-hosts.
|
||||
Transport is stdlib ``urllib`` in a thread executor (no HTTP client deps).
|
||||
The relay wire carries media BY REFERENCE: inbound ``media_urls`` name connector
|
||||
re-hosts (``{connector}/relay/media/{id}``) and an outbound ``send_media`` op names
|
||||
a ``source_url`` the connector resolves back to bytes. ``download(url)`` GETs a
|
||||
re-hosted attachment to a local temp file (vision/file tools consume LOCAL paths,
|
||||
like every native adapter); ``upload(path)`` POSTs local bytes to ``/relay/media``
|
||||
and returns the reference for a subsequent ``send_media`` op. Both present the
|
||||
per-gateway signed bearer the WS upgrade uses; stdlib ``urllib`` in a thread executor.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -45,7 +38,7 @@ _MEDIA_USER_AGENT = "HermesAgent-Relay/1.0 (+https://github.com/NousResearch/her
|
||||
|
||||
|
||||
def media_base_url(relay_dial_url: str) -> str:
|
||||
"""Map the ``ws(s)://…/relay`` dial URL to the ``http(s)://…`` base (same derivation as ``_provision_url``)."""
|
||||
"""Map the ``ws(s)://…/relay`` dial URL to the ``http(s)://…`` connector base."""
|
||||
raw = (relay_dial_url or "").strip().rstrip("/")
|
||||
if raw.startswith("ws://"):
|
||||
raw = "http://" + raw[len("ws://") :]
|
||||
@@ -59,12 +52,7 @@ def media_base_url(relay_dial_url: str) -> str:
|
||||
class RelayMediaClient:
|
||||
"""Authenticated client for the connector's ``/relay/media`` routes."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
base_url: str,
|
||||
gateway_id: Optional[str],
|
||||
secret: Optional[str],
|
||||
) -> None:
|
||||
def __init__(self, base_url: str, gateway_id: Optional[str], secret: Optional[str]) -> None:
|
||||
self._base_url = base_url.rstrip("/")
|
||||
self._gateway_id = gateway_id or ""
|
||||
self._secret = secret or ""
|
||||
@@ -77,16 +65,8 @@ class RelayMediaClient:
|
||||
def _bearer(self) -> str:
|
||||
return make_upgrade_token(self._gateway_id, self._secret)
|
||||
|
||||
def is_relay_media_url(self, url: str) -> bool:
|
||||
"""Is ``url`` a connector re-host reference (needs our bearer to GET)?"""
|
||||
return "/relay/media/" in (url or "")
|
||||
|
||||
async def upload(
|
||||
self,
|
||||
file_path: str,
|
||||
*,
|
||||
mime: Optional[str] = None,
|
||||
filename: Optional[str] = None,
|
||||
self, file_path: str, *, mime: Optional[str] = None, filename: Optional[str] = None
|
||||
) -> Optional[str]:
|
||||
"""POST local file bytes to ``/relay/media``; return the reference URL or None on any failure."""
|
||||
if not self.enabled:
|
||||
@@ -99,16 +79,11 @@ class RelayMediaClient:
|
||||
return None
|
||||
if not data or len(data) > MEDIA_MAX_BYTES:
|
||||
logger.warning(
|
||||
"relay media upload: %s size %d outside (0, %d]",
|
||||
file_path,
|
||||
len(data),
|
||||
MEDIA_MAX_BYTES,
|
||||
"relay media upload: %s size %d outside (0, %d]", file_path, len(data), MEDIA_MAX_BYTES
|
||||
)
|
||||
return None
|
||||
content_type = (
|
||||
mime
|
||||
or mimetypes.guess_type(filename or path.name)[0]
|
||||
or "application/octet-stream"
|
||||
mime or mimetypes.guess_type(filename or path.name)[0] or "application/octet-stream"
|
||||
)
|
||||
headers = {
|
||||
"User-Agent": _MEDIA_USER_AGENT,
|
||||
@@ -122,11 +97,8 @@ class RelayMediaClient:
|
||||
req = urllib.request.Request(url, data=data, headers=headers, method="POST")
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=_REQUEST_TIMEOUT_S) as resp:
|
||||
body = json.loads(resp.read().decode("utf-8"))
|
||||
media_id = body.get("id")
|
||||
if not media_id:
|
||||
return None
|
||||
return f"{self._base_url}/relay/media/{media_id}"
|
||||
media_id = json.loads(resp.read().decode("utf-8")).get("id")
|
||||
return f"{self._base_url}/relay/media/{media_id}" if media_id else None
|
||||
except (urllib.error.URLError, ValueError, OSError) as exc:
|
||||
logger.warning("relay media upload failed: %s", exc)
|
||||
return None
|
||||
@@ -141,7 +113,7 @@ class RelayMediaClient:
|
||||
"""
|
||||
if not url:
|
||||
return None
|
||||
needs_auth = self.is_relay_media_url(url)
|
||||
needs_auth = "/relay/media/" in url
|
||||
if needs_auth and not self.enabled:
|
||||
return None
|
||||
headers = {"User-Agent": _MEDIA_USER_AGENT}
|
||||
@@ -152,8 +124,7 @@ class RelayMediaClient:
|
||||
req = urllib.request.Request(url, headers=headers)
|
||||
try:
|
||||
with urllib.request.urlopen(req, timeout=_REQUEST_TIMEOUT_S) as resp:
|
||||
length = int(resp.headers.get("Content-Length") or 0)
|
||||
if length > MEDIA_MAX_BYTES:
|
||||
if int(resp.headers.get("Content-Length") or 0) > MEDIA_MAX_BYTES:
|
||||
logger.warning("relay media download too large: %s", url)
|
||||
return None
|
||||
data = resp.read(MEDIA_MAX_BYTES + 1)
|
||||
|
||||
+22
-37
@@ -1,15 +1,10 @@
|
||||
"""Relay transport protocol — the gateway<->connector wire contract. EXPERIMENTAL.
|
||||
|
||||
The ``RelayAdapter`` (gateway side) delegates all wire I/O to a ``RelayTransport``.
|
||||
The gateway dials OUT to the connector, so a production transport is a WebSocket
|
||||
client (``ws_transport.py``); in tests it is an in-memory stub
|
||||
(``tests/gateway/relay/stub_connector.py``). This module defines the protocol
|
||||
surface only: lifecycle (connect/disconnect), handshake (the advertised
|
||||
``CapabilityDescriptor``), inbound (``set_inbound_handler``) and outbound
|
||||
(``send_outbound`` / ``get_chat_info`` / ``send_interrupt``).
|
||||
|
||||
EXPERIMENTAL: may change without a deprecation cycle until >=2 Class-1 platforms
|
||||
validate it. See docs/relay-connector-contract.md.
|
||||
The ``RelayAdapter`` delegates all wire I/O to a ``RelayTransport``. The gateway
|
||||
dials OUT to the connector, so production is a WebSocket client (``ws_transport.py``)
|
||||
and tests use an in-memory stub (``tests/gateway/relay/stub_connector.py``). This
|
||||
module defines the protocol surface only. May change without a deprecation cycle
|
||||
until >=2 Class-1 platforms validate it. See docs/relay-connector-contract.md.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -38,7 +33,6 @@ class RelayTransport(Protocol):
|
||||
...
|
||||
|
||||
async def disconnect(self) -> None:
|
||||
"""Close the connection."""
|
||||
...
|
||||
|
||||
async def handshake(self) -> CapabilityDescriptor:
|
||||
@@ -46,15 +40,14 @@ class RelayTransport(Protocol):
|
||||
...
|
||||
|
||||
def set_inbound_handler(self, handler: InboundHandler) -> None:
|
||||
"""Register the callback invoked with each inbound MessageEvent."""
|
||||
...
|
||||
|
||||
def set_passthrough_handler(self, handler: "PassthroughHandler") -> None:
|
||||
"""Register the callback for each forwarded passthrough request (§5.1).
|
||||
|
||||
The connector answers the provider's edge ACK itself, then forwards the
|
||||
real request over this same outbound socket (a hosted gateway has no
|
||||
public inbound port). Optional on a transport (a stub may not implement it).
|
||||
The connector answers the provider's edge ACK itself, then forwards the real
|
||||
request over this same outbound socket (a hosted gateway has no public
|
||||
inbound port). Optional on a transport (a stub may not implement it).
|
||||
"""
|
||||
...
|
||||
|
||||
@@ -75,22 +68,18 @@ class RelayTransport(Protocol):
|
||||
...
|
||||
|
||||
async def send_interrupt(self, session_key: str, reason: Optional[str] = None) -> None:
|
||||
"""Route a mid-turn /stop to the connector for ``session_key``.
|
||||
|
||||
The connector forwards it down the socket owned by the gateway instance
|
||||
running that session. This is the OUTBOUND direction; the actual task
|
||||
cancellation happens when the connector echoes an interrupt inbound.
|
||||
"""
|
||||
"""Route a mid-turn /stop to the connector for ``session_key`` (OUTBOUND
|
||||
direction; the actual cancellation happens when the connector echoes an
|
||||
interrupt inbound down the socket owning that session)."""
|
||||
...
|
||||
|
||||
async def go_idle(self, timeout_s: float = 10.0) -> bool:
|
||||
"""Ask the connector to flip this instance to buffered-only (§5.3).
|
||||
|
||||
Sends ``going_idle`` and awaits the connector-authoritative
|
||||
``going_idle_ack`` (live delivery stopped; inbound now buffers for replay
|
||||
on reconnect). Returns True on ack, False on timeout / not-connected —
|
||||
the caller closes regardless. Optional on a transport. Emitted as part of
|
||||
the gateway's EXISTING drain transition -- not a new idle path.
|
||||
Sends ``going_idle`` and awaits the connector-authoritative ``going_idle_ack``
|
||||
(live delivery stopped; inbound now buffers for replay on reconnect). True on
|
||||
ack, False on timeout / not-connected — the caller closes regardless.
|
||||
Optional on a transport; part of the gateway's EXISTING drain transition.
|
||||
"""
|
||||
...
|
||||
|
||||
@@ -99,16 +88,12 @@ class RelayTransport(Protocol):
|
||||
) -> Dict[str, Any]:
|
||||
"""Act on a shared-identity capability bound to a session (A2 outbound).
|
||||
|
||||
Some platforms hand the connector a credential acting on the SHARED bot
|
||||
identity (e.g. a Discord interaction follow-up token). Under A2 it NEVER
|
||||
reaches the gateway: the connector binds it in its capability vault keyed
|
||||
by session, and the gateway issues a SEMANTIC action against that session.
|
||||
|
||||
Action dict: ``op == "follow_up"``, ``session_key``, ``kind`` (e.g.
|
||||
``"discord.interaction_token"``), ``content``, optional ``metadata``.
|
||||
The connector resolves the capability, enforces the tenant match, and
|
||||
egresses. Returns ``{success, message_id?, error?}``; ``success`` is False
|
||||
when the capability is absent/expired or the tenant mismatches — nothing to
|
||||
retry with (by design: a leaked gateway holds zero capability material).
|
||||
A credential acting on the SHARED bot identity (e.g. a Discord interaction
|
||||
follow-up token) NEVER reaches the gateway: the connector vaults it keyed by
|
||||
session and the gateway issues a SEMANTIC action (``op == "follow_up"``,
|
||||
``session_key``, ``kind`` e.g. ``"discord.interaction_token"``, ``content``,
|
||||
optional ``metadata``). Returns ``{success, message_id?, error?}``; ``success``
|
||||
is False when the capability is absent/expired or the tenant mismatches —
|
||||
nothing to retry with (a leaked gateway holds zero capability material).
|
||||
"""
|
||||
...
|
||||
|
||||
Reference in New Issue
Block a user