# Conflicts: # AGENTS.md # acp_adapter/edit_approval.py # acp_adapter/server.py # agent/agent_init.py # agent/anthropic_adapter.py # agent/anthropic_credentials.py # agent/auxiliary_client.py # agent/azure_identity_adapter.py # agent/bedrock_adapter.py # agent/browser_registry.py # agent/chat_completion_helpers.py # agent/coding_context.py # agent/context_references.py # agent/conversation_loop.py # agent/copilot_acp_client.py # agent/credits_tracker.py # agent/curator.py # agent/curator_backup.py # agent/deadline.py # agent/display.py # agent/errors.py # agent/estop.py # agent/i18n.py # agent/image_gen_registry.py # agent/image_routing.py # agent/learning_graph.py # agent/learning_mutations.py # agent/lsp/servers.py # agent/model_metadata.py # agent/models_dev.py # agent/monitoring/gateway_health_export.py # agent/monitoring/otlp_exporter.py # agent/pet/store.py # agent/process_bootstrap.py # agent/prompt_builder.py # agent/proxy_sources/iron_proxy.py # agent/secret_sources/_cache.py # agent/secret_sources/bitwarden.py # agent/secret_sources/registry.py # agent/shell_hooks.py # agent/skill_bundles.py # agent/skill_commands.py # agent/skill_utils.py # agent/ssl_guard.py # agent/ssl_verify.py # agent/system_prompt.py # agent/terminal_env_registry.py # agent/trace_upload.py # agent/transcription_registry.py # agent/tts_registry.py # agent/verify/environment.py # agent/vertex_adapter.py # agent/video_gen_registry.py # agent/web_search_registry.py # cli.py # cron/jobs.py # cron/scheduler.py # gateway/agent_cache_pressure.py # gateway/cgroup_cleanup.py # gateway/channel_directory.py # gateway/config.py # gateway/control_socket.py # gateway/dead_targets.py # gateway/drain_control.py # gateway/hooks.py # gateway/kanban_watchers.py # gateway/lifecycle_ledger.py # gateway/mirror.py # gateway/pairing.py # gateway/platform_registry.py # gateway/platforms/helpers.py # gateway/platforms/weixin.py # gateway/readiness.py # gateway/restart_loop_guard.py # gateway/rich_sent_store.py # gateway/run.py # gateway/session.py # gateway/shutdown_flush.py # gateway/shutdown_forensics.py # gateway/slash_commands.py # gateway/status.py # gateway/sticker_cache.py # gateway/whatsapp_identity.py # hermes_bootstrap.py # hermes_cli/_early_recovery.py # hermes_cli/_install_repair.py # hermes_cli/_startup_fast.py # hermes_cli/_subprocess_compat.py # hermes_cli/agent_plugins.py # hermes_cli/auth.py # hermes_cli/backup.py # hermes_cli/banner.py # hermes_cli/browser_connect.py # hermes_cli/build_info.py # hermes_cli/cli_agent_setup_mixin.py # hermes_cli/cli_commands_mixin.py # hermes_cli/codex_models.py # hermes_cli/config.py # hermes_cli/config_defaults.py # hermes_cli/config_migrations.py # hermes_cli/container_boot.py # hermes_cli/dashboard_auth/registry.py # hermes_cli/debug.py # hermes_cli/dep_ensure.py # hermes_cli/doctor.py # hermes_cli/doctor_live.py # hermes_cli/dump.py # hermes_cli/env_loader.py # hermes_cli/foreign_sessions.py # hermes_cli/gateway.py # hermes_cli/gateway_windows.py # hermes_cli/gui_uninstall.py # hermes_cli/image_provenance.py # hermes_cli/install_identity.py # hermes_cli/kanban.py # hermes_cli/kanban_db.py # hermes_cli/linux_desktop_entry.py # hermes_cli/local_runtime/binaries.py # hermes_cli/local_runtime/endpoint.py # hermes_cli/local_runtime/growth.py # hermes_cli/local_runtime/supervisor.py # hermes_cli/logs.py # hermes_cli/macos_tcc_anchor.py # hermes_cli/main.py # hermes_cli/memory_setup.py # hermes_cli/model_catalog.py # hermes_cli/models.py # hermes_cli/nous_subscription.py # hermes_cli/npm_engine.py # hermes_cli/plugin_index.py # hermes_cli/plugins.py # hermes_cli/plugins_cmd.py # hermes_cli/profile_distribution.py # hermes_cli/profiles.py # hermes_cli/prompt_size.py # hermes_cli/psutil_android.py # hermes_cli/runtime_repair.py # hermes_cli/security_advisories.py # hermes_cli/security_audit.py # hermes_cli/security_audit_startup.py # hermes_cli/service_manager.py # hermes_cli/session_export_md.py # hermes_cli/setup.py # hermes_cli/skills_hub.py # hermes_cli/slack_cli.py # hermes_cli/status.py # hermes_cli/subcommands/gateway.py # hermes_cli/subcommands/uninstall.py # hermes_cli/tools_config.py # hermes_cli/uninstall.py # hermes_cli/update_cmd.py # hermes_cli/update_contract.py # hermes_cli/update_inventory.py # hermes_cli/update_lock.py # hermes_cli/update_receipt.py # hermes_cli/urllib_security.py # hermes_cli/web_routers/local_models.py # hermes_cli/web_routers/profiles.py # hermes_cli/web_routers/skills.py # hermes_cli/web_server.py # hermes_constants.py # hermes_state.py # plugins/disk-cleanup/__init__.py # plugins/disk-cleanup/disk_cleanup.py # plugins/google_meet/node/registry.py # plugins/google_meet/node/server.py # plugins/google_meet/process_manager.py # plugins/google_meet/realtime/openai_client.py # plugins/hermes-achievements/dashboard/plugin_api.py # plugins/memory/hindsight/__init__.py # plugins/memory/honcho/__init__.py # plugins/memory/honcho/cli.py # plugins/memory/honcho/client.py # plugins/memory/honcho/oauth.py # plugins/memory/honcho/session.py # plugins/memory/mem0/__init__.py # plugins/memory/mem0/_setup.py # plugins/memory/openviking/__init__.py # plugins/memory/retaindb/__init__.py # plugins/memory/supermemory/__init__.py # plugins/platforms/a2a/protocol.py # plugins/platforms/dingtalk/adapter.py # plugins/platforms/discord/adapter.py # plugins/platforms/feishu/adapter.py # plugins/platforms/google_chat/adapter.py # plugins/platforms/matrix/adapter.py # plugins/platforms/photon/adapter.py # plugins/platforms/photon/auth.py # plugins/platforms/photon/cli.py # plugins/platforms/slack/adapter.py # plugins/platforms/teams/adapter.py # plugins/platforms/telegram/adapter.py # plugins/platforms/wecom/callback_adapter.py # plugins/platforms/whatsapp/adapter.py # plugins/teams_pipeline/store.py # plugins/video_gen/fal/__init__.py # plugins/web/ddgs/provider.py # plugins/web/exa/provider.py # plugins/web/firecrawl/provider.py # plugins/web/parallel/provider.py # tests/agent/test_ssl_ca_guard.py # tests/hermes_cli/test_certifi_repair.py # tests/hermes_cli/test_cmd_update.py # tests/hermes_cli/test_cmd_update_apt.py # tests/hermes_cli/test_dashboard_unified_launch.py # tests/hermes_cli/test_dep_ensure.py # tests/hermes_cli/test_doctor.py # tests/hermes_cli/test_doctor_live.py # tests/hermes_cli/test_gui_command.py # tests/hermes_cli/test_kanban_boards.py # tests/hermes_cli/test_kanban_db.py # tests/hermes_cli/test_lazy_refresh_venv_repair.py # tests/hermes_cli/test_memory_setup_provider_arg.py # tests/hermes_cli/test_nous_subscription.py # tests/hermes_cli/test_pip_install_detection.py # tests/hermes_cli/test_profile_export_credentials.py # tests/hermes_cli/test_psutil_android_extract.py # tests/hermes_cli/test_status.py # tests/hermes_cli/test_tui_npm_install.py # tests/hermes_cli/test_update_fleet_restart_pending.py # tests/hermes_cli/test_update_head_moved_gate.py # tests/hermes_cli/test_update_interrupted_recovery.py # tests/hermes_cli/test_web_server.py # tests/hermes_cli/test_web_ui_build.py # tests/test_hermes_logging.py # tests/test_managed_runtime_resolution.py # tests/tools/test_browser_chromium_autoinstall.py # tests/tools/test_browser_chromium_check.py # tests/tools/test_browser_homebrew_paths.py # tests/tools/test_browser_lightpanda.py # tests/tools/test_browser_npx_warmup.py # tests/tools/test_browser_open_timeout.py # tests/tools/test_browser_orphan_reaper.py # tests/tools/test_browser_real_profile.py # tests/tools/test_browser_suspect_recycle.py # tests/tools/test_find_shell.py # tests/tools/test_local_env_blocklist.py # tests/tools/test_macos_protected_search.py # tests/tui_gateway/test_compute_host.py # tools/approval.py # tools/blueprints.py # tools/bot_mode_dm.py # tools/bot_mode_probe.py # tools/bot_relay.py # tools/browser_tool.py # tools/browser_use_cli.py # tools/checkpoint_manager.py # tools/code_execution_tool.py # tools/code_kernel.py # tools/computer_use/cua_backend.py # tools/cronjob_tools.py # tools/discord_tool.py # tools/environments/base.py # tools/environments/daytona.py # tools/environments/local.py # tools/environments/modal.py # tools/environments/vercel_sandbox.py # tools/fal_common.py # tools/file_operations.py # tools/lazy_deps.py # tools/mcp_tool.py # tools/neutts_synth.py # tools/process_registry.py # tools/read_extract.py # tools/registry.py # tools/skill_ledger.py # tools/skill_linter.py # tools/skill_manager_tool.py # tools/skill_usage.py # tools/skills_ast_audit.py # tools/skills_guard.py # tools/skills_hub.py # tools/skills_sync.py # tools/skills_sync_client.py # tools/skills_tool.py # tools/terminal_scope.py # tools/terminal_tool.py # tools/tirith_security.py # tools/transcription_tools.py # tools/tts_tool.py # tools/vision_tools.py # tools/voice_mode.py # tools/wake_word.py # tools/web_result_cache.py # tools/website_policy.py # tools/working_diff.py # tools/write_approval.py # tui_gateway/entry.py # tui_gateway/methods_tools.py # tui_gateway/server.py
343 lines
16 KiB
Python
343 lines
16 KiB
Python
"""Gateway control socket — the gateway-owned local coordination surface: a local-only socket answering
|
|
versioned JSON verbs (``identify``, ``status``). A connectable socket with a well-formed ``identify``
|
|
answer IS liveness — no PID-reuse heuristics. Never a TCP port: filesystem/pipe ACLs are the auth
|
|
boundary. POSIX: ``$HERMES_HOME/gateway.sock`` (or a temp-dir socket + ``gateway.sock.path`` pointer
|
|
file when the home path exceeds ``sun_path``); Windows: named pipe ``\\\\.\\pipe\\hermes-gateway-<hash>``.
|
|
Wire contract: ONE request per connection — one JSON line in, one out, then the server closes.
|
|
Consumers PREFER the socket and fall back to the state-file/scan layer when it doesn't answer.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import hashlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import socket
|
|
import sys
|
|
import tempfile
|
|
import time
|
|
from pathlib import Path
|
|
from typing import Any, Callable, Optional
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
CONTROL_PROTOCOL_VERSION = 1
|
|
_SOCKET_FILENAME = "gateway.sock"
|
|
_POINTER_FILENAME = "gateway.sock.path"
|
|
_IS_WINDOWS = sys.platform == "win32"
|
|
_MAX_UNIX_PATH = 100 # sun_path limit is 104 on macOS/BSD, 108 on Linux; margin for the NUL
|
|
# Single-line JSON in/out; bounded so a misbehaving peer can't balloon memory.
|
|
_MAX_REQUEST_BYTES = 64 * 1024
|
|
_MAX_RESPONSE_BYTES = 512 * 1024
|
|
_DEFAULT_CLIENT_TIMEOUT = 2.0
|
|
|
|
|
|
def _home_hash(home: Path) -> str:
|
|
return hashlib.sha256(os.path.normcase(str(Path(home).expanduser().resolve(strict=False))).encode("utf-8")).hexdigest()[:16]
|
|
|
|
|
|
def windows_pipe_name(home: Path) -> str:
|
|
"""Per-HERMES_HOME named pipe path (Windows transport)."""
|
|
return rf"\\.\pipe\hermes-gateway-{_home_hash(home)}"
|
|
|
|
|
|
def _fits_sun_path(path: Path) -> bool:
|
|
return len(str(path).encode("utf-8")) <= _MAX_UNIX_PATH
|
|
|
|
|
|
def _fallback_socket_path(home: Path) -> Path:
|
|
"""Short temp-dir path for homes whose direct socket path exceeds sun_path: ``tempfile.gettempdir()``
|
|
then ``/tmp`` (POSIX); if nothing fits the tempdir candidate is returned anyway — bind fails
|
|
non-fatally and consumers use the scan layer."""
|
|
name = f"hermes-gw-{_home_hash(home)}.sock"
|
|
candidates = [Path(tempfile.gettempdir()) / name] + ([] if _IS_WINDOWS else [Path("/tmp") / name])
|
|
return next((c for c in candidates if _fits_sun_path(c)), candidates[0])
|
|
|
|
|
|
def resolve_server_socket_path(home: Path) -> tuple[Path, Optional[Path]]:
|
|
"""Return ``(bind_path, pointer_file)``; pointer_file is set only for the temp-dir fallback."""
|
|
direct = Path(home) / _SOCKET_FILENAME
|
|
return (direct, None) if _fits_sun_path(direct) else (_fallback_socket_path(home), Path(home) / _POINTER_FILENAME)
|
|
|
|
|
|
def resolve_client_socket_path(home: Path) -> Optional[Path]:
|
|
"""Where a client should connect for ``home``, or None when nothing exists."""
|
|
direct = Path(home) / _SOCKET_FILENAME
|
|
if direct.exists():
|
|
return direct
|
|
with contextlib.suppress(OSError):
|
|
pointer = Path(home) / _POINTER_FILENAME
|
|
target = pointer.read_text(encoding="utf-8-sig").strip() if pointer.is_file() else ""
|
|
if target and Path(target).exists():
|
|
return Path(target)
|
|
return None
|
|
|
|
|
|
def _detect_supervisor() -> str:
|
|
"""Supervisor kind for THIS process from its own launch env (not inferred outside-in).
|
|
|
|
Unlike the outside-in `_detect_supervisor_for_pid` scan, this answers from the process's own launch
|
|
context — which is exactly the provenance the 92091 design wants declared rather than inferred. See
|
|
#92091.
|
|
"""
|
|
env = os.environ
|
|
if env.get("INVOCATION_ID"):
|
|
return "systemd"
|
|
if sys.platform == "darwin" and (env.get("XPC_SERVICE_NAME", "").startswith("ai.hermes")
|
|
or env.get("LAUNCHD_SOCKET")):
|
|
return "launchd"
|
|
if env.get("HERMES_DESKTOP_MANAGED"):
|
|
return "desktop"
|
|
return "external" if "--external-supervisor" in sys.argv else "manual"
|
|
|
|
|
|
def build_identify_payload() -> dict[str, Any]:
|
|
"""Default ``identify`` answer, built from gateway.status primitives."""
|
|
from gateway.status import _build_pid_record, _get_code_identity_fields, _profile_label_for_home, read_runtime_status
|
|
record = _build_pid_record()
|
|
payload: dict[str, Any] = {
|
|
"protocol": CONTROL_PROTOCOL_VERSION,
|
|
**{k: record.get(k) for k in ("kind", "pid", "start_time", "hermes_home")},
|
|
"profile": _profile_label_for_home(record.get("hermes_home") or ""),
|
|
"supervisor": _detect_supervisor(), **_get_code_identity_fields()}
|
|
with contextlib.suppress(Exception):
|
|
# served_profiles (multiplex mode) is stamped into runtime status by the runner.
|
|
served = (read_runtime_status() or {}).get("served_profiles")
|
|
if isinstance(served, list) and served:
|
|
payload["served_profiles"] = served
|
|
return payload
|
|
|
|
|
|
def build_status_payload() -> dict[str, Any]:
|
|
"""Default ``status`` answer — current runtime status, answered live."""
|
|
from gateway.status import read_runtime_status
|
|
return {**(read_runtime_status() or {}), "protocol": CONTROL_PROTOCOL_VERSION,
|
|
"answered_at": time.time(), "answering_pid": os.getpid()}
|
|
|
|
|
|
class GatewayControlServer:
|
|
"""Gateway-owned control socket server (identify/status, v1): ``start()`` after the PID-file claim,
|
|
``stop()`` on shutdown. All failures are non-fatal — the gateway never refuses to serve messaging
|
|
because its control socket couldn't bind; consumers fall back to the scan layer."""
|
|
|
|
def __init__(self, home: Optional[Path] = None, *,
|
|
verb_handlers: Optional[dict[str, Callable[[], dict[str, Any]]]] = None) -> None:
|
|
if home is None:
|
|
from gateway.status import _get_process_hermes_home
|
|
home = _get_process_hermes_home()
|
|
self._home = Path(home)
|
|
self._server: Optional[asyncio.AbstractServer] = None
|
|
self._pipe_server: Any = None # Windows proactor pipe server
|
|
self._bind_path: Optional[Path] = None
|
|
self._pointer_file: Optional[Path] = None
|
|
self._handlers: dict[str, Callable[[], dict[str, Any]]] = {
|
|
"identify": build_identify_payload, "status": build_status_payload, **(verb_handlers or {})}
|
|
|
|
async def start(self) -> bool:
|
|
"""Bind and start serving. Returns True on success, False otherwise."""
|
|
try:
|
|
return await (self._start_windows() if _IS_WINDOWS else self._start_posix())
|
|
except Exception as exc:
|
|
logger.warning("Gateway control socket failed to start (non-fatal): %s", exc)
|
|
return False
|
|
|
|
async def _start_posix(self) -> bool:
|
|
bind_path, pointer_file = resolve_server_socket_path(self._home)
|
|
# We only get here after winning the PID-file O_EXCL race, so any existing
|
|
# file is stale or a collision — never a live sibling.
|
|
with contextlib.suppress(OSError):
|
|
if bind_path.exists():
|
|
bind_path.unlink()
|
|
# Restrictive umask so the socket is never world-connectable, even for the instant before chmod.
|
|
old_umask = os.umask(0o177)
|
|
try:
|
|
self._server = await asyncio.start_unix_server(self._handle_connection, path=str(bind_path))
|
|
finally:
|
|
os.umask(old_umask)
|
|
with contextlib.suppress(OSError):
|
|
os.chmod(bind_path, 0o600)
|
|
self._bind_path = bind_path
|
|
if pointer_file is not None:
|
|
pointer_file.write_text(str(bind_path), encoding="utf-8")
|
|
self._pointer_file = pointer_file
|
|
logger.info("Gateway control socket listening at %s", bind_path)
|
|
return True
|
|
|
|
async def _start_windows(self) -> bool:
|
|
loop = asyncio.get_running_loop()
|
|
start_serving_pipe = getattr(loop, "start_serving_pipe", None)
|
|
if start_serving_pipe is None:
|
|
logger.debug("Event loop %s has no start_serving_pipe — control socket "
|
|
"disabled (selector loop on Windows).", type(loop).__name__)
|
|
return False
|
|
pipe_name = windows_pipe_name(self._home)
|
|
servers = await start_serving_pipe(lambda: _PipeControlProtocol(self), pipe_name)
|
|
self._pipe_server = servers[0] if servers else None
|
|
logger.info("Gateway control pipe listening at %s", pipe_name)
|
|
return self._pipe_server is not None
|
|
|
|
async def stop(self) -> None:
|
|
"""Stop serving and remove the socket/pointer files."""
|
|
if self._server is not None:
|
|
self._server.close()
|
|
with contextlib.suppress(Exception):
|
|
await self._server.wait_closed()
|
|
if self._pipe_server is not None:
|
|
with contextlib.suppress(Exception):
|
|
self._pipe_server.close()
|
|
self._server = self._pipe_server = None
|
|
self.cleanup_files()
|
|
|
|
def cleanup_files(self) -> None:
|
|
"""Best-effort removal of socket + pointer files (atexit-safe)."""
|
|
for path in filter(None, (self._bind_path, self._pointer_file)):
|
|
with contextlib.suppress(OSError):
|
|
path.unlink(missing_ok=True)
|
|
|
|
def handle_request_line(self, raw: bytes) -> bytes:
|
|
"""One JSON request line -> one JSON response line. Never raises (shared by POSIX + pipe)."""
|
|
request_id: Any = None
|
|
try:
|
|
request = json.loads(raw.decode("utf-8"))
|
|
if not isinstance(request, dict):
|
|
raise ValueError("request must be a JSON object")
|
|
request_id, verb = request.get("id"), request.get("verb")
|
|
handler = self._handlers.get(verb) if isinstance(verb, str) else None
|
|
if handler is None:
|
|
response: dict[str, Any] = {"ok": False, "error": f"unknown verb: {verb!r}",
|
|
"protocol": CONTROL_PROTOCOL_VERSION, "supported_verbs": sorted(self._handlers)}
|
|
else:
|
|
response = {"ok": True, "protocol": CONTROL_PROTOCOL_VERSION, "result": handler()}
|
|
except Exception as exc:
|
|
response = {"ok": False, "error": f"{type(exc).__name__}: {exc}", "protocol": CONTROL_PROTOCOL_VERSION}
|
|
if request_id is not None:
|
|
response["id"] = request_id
|
|
try:
|
|
encoded = json.dumps(response, default=str).encode("utf-8")
|
|
except Exception:
|
|
encoded = b'{"ok": false, "error": "response serialization failed"}'
|
|
if len(encoded) > _MAX_RESPONSE_BYTES:
|
|
encoded = b'{"ok": false, "error": "response too large"}'
|
|
return encoded + b"\n"
|
|
|
|
async def _handle_connection(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
|
|
try:
|
|
raw = await asyncio.wait_for(reader.readline(), timeout=_DEFAULT_CLIENT_TIMEOUT)
|
|
if not raw or len(raw) > _MAX_REQUEST_BYTES:
|
|
return
|
|
# Handlers read disk; keep that off the loop that drives every platform
|
|
# adapter so a fast-polling consumer can't stall heartbeats.
|
|
response = await asyncio.get_running_loop().run_in_executor(
|
|
None, self.handle_request_line, raw.rstrip(b"\n"))
|
|
writer.write(response)
|
|
await writer.drain()
|
|
except (asyncio.TimeoutError, ConnectionError, OSError):
|
|
pass
|
|
except Exception:
|
|
logger.debug("Control socket connection handler error", exc_info=True)
|
|
finally:
|
|
with contextlib.suppress(Exception):
|
|
writer.close()
|
|
|
|
|
|
class _PipeControlProtocol(asyncio.Protocol):
|
|
"""One-shot request/response protocol for the Windows named pipe."""
|
|
def __init__(self, server: GatewayControlServer) -> None:
|
|
self._server = server
|
|
self._transport: Any = None
|
|
self._buffer = bytearray()
|
|
|
|
def connection_made(self, transport) -> None: # pragma: no cover - windows
|
|
self._transport = transport
|
|
|
|
def data_received(self, data: bytes) -> None: # pragma: no cover - windows
|
|
self._buffer.extend(data)
|
|
if len(self._buffer) > _MAX_REQUEST_BYTES:
|
|
self._transport.close()
|
|
elif b"\n" in self._buffer:
|
|
try:
|
|
self._transport.write(self._server.handle_request_line(bytes(self._buffer).partition(b"\n")[0]))
|
|
finally:
|
|
self._transport.close()
|
|
|
|
|
|
def query_gateway_control(home: Path, verb: str, *, timeout: float = _DEFAULT_CLIENT_TIMEOUT) -> Optional[dict[str, Any]]:
|
|
"""Ask the gateway serving ``home`` a control verb; returns its ``result`` payload. Any failure (no/stale
|
|
socket, timeout, malformed answer, ``ok: false``) returns None so callers fall back to the scan layer.
|
|
Never raises."""
|
|
request = json.dumps({"verb": verb, "id": 1, "protocol": CONTROL_PROTOCOL_VERSION}).encode("utf-8") + b"\n"
|
|
query = _query_windows_pipe if _IS_WINDOWS else _query_unix_socket
|
|
try:
|
|
raw = query(Path(home), request, timeout)
|
|
response = json.loads(raw.decode("utf-8")) if raw else None
|
|
except Exception:
|
|
return None
|
|
result = response.get("result") if isinstance(response, dict) and response.get("ok") is True else None
|
|
return result if isinstance(result, dict) else None
|
|
|
|
|
|
def _read_response_line(read: Callable[[], bytes], deadline: float) -> Optional[bytes]:
|
|
"""Read chunks until a newline, EOF, deadline, or the size cap (-> None)."""
|
|
chunks: list[bytes] = []
|
|
while time.monotonic() < deadline:
|
|
chunk = read()
|
|
chunks.append(chunk)
|
|
if not chunk or b"\n" in chunk:
|
|
break
|
|
if sum(len(c) for c in chunks) > _MAX_RESPONSE_BYTES:
|
|
return None
|
|
return b"".join(chunks).partition(b"\n")[0] or None
|
|
|
|
|
|
def _query_unix_socket(home: Path, request: bytes, timeout: float) -> Optional[bytes]:
|
|
path = resolve_client_socket_path(home)
|
|
if path is None:
|
|
return None
|
|
# OSError covers ConnectionRefusedError / FileNotFoundError on connect and socket.timeout on read.
|
|
with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as sock, contextlib.suppress(OSError):
|
|
sock.settimeout(timeout)
|
|
sock.connect(str(path))
|
|
sock.sendall(request)
|
|
return _read_response_line(lambda: sock.recv(65536), time.monotonic() + timeout)
|
|
return None
|
|
|
|
|
|
def _query_windows_pipe(home: Path, request: bytes, timeout: float) -> Optional[bytes]: # pragma: no cover - wine2e lane
|
|
pipe_name = windows_pipe_name(home)
|
|
deadline = time.monotonic() + timeout
|
|
handle = None
|
|
while handle is None:
|
|
try:
|
|
handle = open(pipe_name, "r+b", buffering=0)
|
|
except FileNotFoundError:
|
|
return None
|
|
except OSError:
|
|
# Pipe busy (another client mid-handshake) — brief retry window.
|
|
if time.monotonic() >= deadline:
|
|
return None
|
|
time.sleep(0.05)
|
|
try:
|
|
handle.write(request)
|
|
return _read_response_line(lambda: handle.read(65536), deadline)
|
|
finally:
|
|
with contextlib.suppress(Exception):
|
|
handle.close()
|
|
|
|
|
|
def identify_gateway(home: Path, *, timeout: float = _DEFAULT_CLIENT_TIMEOUT) -> Optional[dict[str, Any]]:
|
|
"""Convenience wrapper: ``identify`` the gateway serving ``home``."""
|
|
return query_gateway_control(home, "identify", timeout=timeout)
|
|
|
|
|
|
def pause_gateway_for_update(home: Path, *, timeout: float = _DEFAULT_CLIENT_TIMEOUT) -> Optional[dict[str, Any]]:
|
|
"""Ask the gateway serving ``home`` to drain and exit for an update. Returns the ACK ``{"pausing",
|
|
"already_stopping", "pid", "drain_timeout"}`` or None when no gateway answers (old gateway without
|
|
the verb, no/dead socket) — the caller then uses the legacy signal/tree-kill pause path.
|
|
|
|
Step 2 of the socket migration (#92091).
|
|
"""
|
|
return query_gateway_control(home, "pause-for-update", timeout=timeout)
|