Application-backed MCP servers: live endpoint per attempt, tools follow app liveness, errors and the tool catalog say why (NS-943) (#119050)
* feat: expose interactive and live endpoint facts * feat: hydrate application-backed MCP liveness * feat: export interactive session fact * fix: classify unsupported declared hosts * test(mcp): a runtime file without a token connects without an Authorization header * fix(mcp): server_json liveness fields default to http/token/pid; an invalid declaration degrades to static instead of failing the server task * fix(mcp): a connected server is offerable regardless of the interactive-session rule; the rule only shapes the not-connected sentence * test(mcp): follow main's ToolSearchConfig fields and patch _bump_server_error at its origin module
This commit is contained in:
@@ -2,3 +2,7 @@
|
||||
|
||||
Facts use hardware sources without environment-variable input or subprocesses.
|
||||
"""
|
||||
|
||||
from hermes_platform.host.facts import interactive_session
|
||||
|
||||
__all__ = ["interactive_session"]
|
||||
|
||||
@@ -170,6 +170,34 @@ def cpu_vendor() -> str:
|
||||
return ""
|
||||
|
||||
|
||||
def _windows_interactive_session() -> bool:
|
||||
import ctypes
|
||||
from ctypes import wintypes
|
||||
|
||||
kernel32 = ctypes.WinDLL("kernel32", use_last_error=True)
|
||||
process_id = kernel32.GetCurrentProcessId()
|
||||
session_id = wintypes.DWORD()
|
||||
if not kernel32.ProcessIdToSessionId(process_id, ctypes.byref(session_id)) or session_id.value == 0:
|
||||
return False
|
||||
active_session = kernel32.WTSGetActiveConsoleSessionId()
|
||||
if active_session != 0xFFFFFFFF and active_session == session_id.value:
|
||||
return True
|
||||
return bool(kernel32.GetProcessWindowStation())
|
||||
|
||||
|
||||
def interactive_session() -> bool:
|
||||
"""Return whether this process can reach an interactive user session."""
|
||||
if sys.platform == "win32":
|
||||
try:
|
||||
return _windows_interactive_session()
|
||||
except (AttributeError, OSError, TypeError, ValueError):
|
||||
return False
|
||||
if sys.platform.startswith("linux"):
|
||||
session_id = _read_text("/proc/self/sessionid", 64).strip()
|
||||
return bool(session_id and session_id != "4294967295" and os.path.isdir(f"/run/user/{os.getuid()}"))
|
||||
return True
|
||||
|
||||
|
||||
def clear_caches() -> None:
|
||||
"""Clear every cached host fact."""
|
||||
for fact in (os_family, process_arch, native_arch, cpu_model, cpu_vendor):
|
||||
|
||||
@@ -54,6 +54,15 @@ def _expand(path: str) -> str:
|
||||
return os.path.expandvars(os.path.expanduser(path))
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Endpoint:
|
||||
url: str
|
||||
token: str = ""
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return f"Endpoint(url={self.url!r}, token=<redacted>)"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class AppResolver:
|
||||
definition: AppDef
|
||||
@@ -100,6 +109,19 @@ class AppResolver:
|
||||
|
||||
# ---- probe: fresh, never cached ------------------------------------------------------
|
||||
|
||||
def endpoint(self) -> Endpoint | None:
|
||||
"""Read and validate the current runtime endpoint."""
|
||||
d = self.definition
|
||||
if d.liveness_kind != "server_json":
|
||||
return None
|
||||
session = _read_server_json(_expand(d.liveness_path), d)
|
||||
if session is None or _pid_alive(session.pid).value is not True:
|
||||
return None
|
||||
endpoint = _endpoint_observation(session.url, d.endpoint_path)
|
||||
if endpoint.state is not CheckState.PRESENT or not endpoint.value:
|
||||
return None
|
||||
return Endpoint(endpoint.value, session.token)
|
||||
|
||||
def probe(self, res: Resolution, *, effort: Effort, deadline_s: float = 3.0) -> Probe:
|
||||
d = self.definition
|
||||
nc: Observation = Observation.not_checked()
|
||||
|
||||
@@ -96,3 +96,7 @@ def test_linux_live_cpu_facts_match_cpuinfo(cleared_fact_caches) -> None:
|
||||
|
||||
assert facts.cpu_vendor() in cpuinfo
|
||||
assert facts.cpu_model() in cpuinfo
|
||||
|
||||
|
||||
def test_interactive_session_is_bool() -> None:
|
||||
assert type(facts.interactive_session()) is bool
|
||||
|
||||
@@ -12,7 +12,7 @@ from http.server import BaseHTTPRequestHandler, HTTPServer
|
||||
import pytest
|
||||
|
||||
from hermes_platform.resolver import CheckState, Effort, Probeable, Resolver
|
||||
from hermes_platform.resolver.app import AppDef, AppResolver
|
||||
from hermes_platform.resolver.app import AppDef, AppResolver, Endpoint
|
||||
|
||||
TOKEN = "tok-3e1f9c-unique-fixture-value"
|
||||
|
||||
@@ -164,6 +164,17 @@ def loopback_mcp():
|
||||
srv.shutdown()
|
||||
|
||||
|
||||
def test_endpoint_reads_server_json_every_call(tmp_path):
|
||||
exe = tmp_path / "thing"
|
||||
exe.write_text("", encoding="utf-8")
|
||||
sj = _server_json(tmp_path, url="http://127.0.0.1:1111/ignored")
|
||||
resolver = _resolver(tmp_path, exe, sj)
|
||||
assert resolver.endpoint() == Endpoint("http://127.0.0.1:1111/mcp", TOKEN)
|
||||
sj.write_text(json.dumps({"pid": os.getpid(), "http": "http://127.0.0.1:2222/ignored", "token": "next"}))
|
||||
assert resolver.endpoint() == Endpoint("http://127.0.0.1:2222/mcp", "next")
|
||||
assert TOKEN not in repr(resolver.endpoint())
|
||||
|
||||
|
||||
def test_network_probe_reads_server_json_every_call_and_answers(tmp_path, loopback_mcp):
|
||||
exe = tmp_path / "thing"
|
||||
exe.write_text("", encoding="utf-8")
|
||||
|
||||
@@ -0,0 +1,191 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
from types import SimpleNamespace
|
||||
|
||||
import pytest
|
||||
|
||||
from hermes_platform import declaration
|
||||
from hermes_platform.resolver.availability import Availability
|
||||
from tools.mcp_liveness import describe, parse_liveness
|
||||
|
||||
|
||||
def _decl(tmp_path, *, min_version=None):
|
||||
executable = tmp_path / "example-app"
|
||||
executable.write_text("fixture", encoding="utf-8")
|
||||
raw = {sys.platform: {"presence": "executable", "location": str(executable)}}
|
||||
requires = {"app": True}
|
||||
if min_version is not None:
|
||||
raw[sys.platform]["version"] = {"kind": "plist" if sys.platform == "darwin" else "none"}
|
||||
requires["min_version"] = min_version
|
||||
return declaration.parse_declaration("Example App", raw, requires, where="test")
|
||||
|
||||
|
||||
def test_parse_liveness_contract():
|
||||
assert parse_liveness({"kind": "static"}).kind == "static"
|
||||
assert parse_liveness({"kind": "interactive_session"}).kind == "interactive_session"
|
||||
live = parse_liveness({
|
||||
"kind": "server_json",
|
||||
"path": "/tmp/example.json",
|
||||
"fields": {"url": "endpoint", "token": "secret", "pid": "process"},
|
||||
})
|
||||
assert (live.kind, live.path, live.url_field, live.token_field, live.pid_field) == (
|
||||
"server_json", "/tmp/example.json", "endpoint", "secret", "process"
|
||||
)
|
||||
with pytest.raises(ValueError, match="unknown liveness kind"):
|
||||
parse_liveness({"kind": "unknown"})
|
||||
defaulted = parse_liveness({"kind": "server_json", "path": "/tmp/example.json"})
|
||||
assert (defaulted.url_field, defaulted.token_field, defaulted.pid_field) == ("http", "token", "pid")
|
||||
partial = parse_liveness({"kind": "server_json", "path": "/tmp/example.json", "fields": {"url": "endpoint"}})
|
||||
assert (partial.url_field, partial.token_field, partial.pid_field) == ("endpoint", "token", "pid")
|
||||
with pytest.raises(ValueError, match="may only override"):
|
||||
parse_liveness({"kind": "server_json", "path": "/tmp/example.json", "fields": {"port": "p"}})
|
||||
|
||||
|
||||
def test_invalid_registered_liveness_degrades_to_static(monkeypatch, caplog):
|
||||
import hermes_cli.agent_plugins as agent_plugins
|
||||
from tools.mcp_liveness import liveness_for
|
||||
|
||||
monkeypatch.setattr(agent_plugins, "liveness_for", lambda name: {"kind": "server_json"}, raising=False)
|
||||
caplog.set_level(logging.WARNING)
|
||||
assert liveness_for("example-server").kind == "static"
|
||||
assert any("invalid liveness declaration" in record.getMessage() for record in caplog.records)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("state", "available", "fragment"),
|
||||
[
|
||||
("app_not_running", Availability("available"), "is not running"),
|
||||
("endpoint_unavailable", Availability("available"), "local endpoint is unavailable"),
|
||||
("no_interactive_session", Availability("available"), "interactive desktop session"),
|
||||
("version_too_old", Availability("version_too_old", version="1.2", min_version="2.0"), "version 1.2 is too old"),
|
||||
("missing_app", Availability("missing_app"), "is not installed"),
|
||||
],
|
||||
)
|
||||
def test_describe_has_one_state_sentence_and_one_action(tmp_path, state, available, fragment):
|
||||
sentence = describe(_decl(tmp_path), available, state)
|
||||
assert sentence.startswith("Example App")
|
||||
assert fragment in sentence
|
||||
assert sentence.count("try again") <= 1
|
||||
|
||||
|
||||
def test_live_endpoint_reloads_file_and_registers_token_before_use(tmp_path, monkeypatch, caplog):
|
||||
import hermes_cli.agent_plugins as agent_plugins
|
||||
from agent import redact
|
||||
from tools.mcp_tool_transport import _live_endpoint
|
||||
|
||||
runtime = tmp_path / "server.json"
|
||||
decl = _decl(tmp_path)
|
||||
declaration.register("example-server", decl)
|
||||
raw = {
|
||||
"kind": "server_json",
|
||||
"path": str(runtime),
|
||||
"fields": {"url": "http", "token": "token", "pid": "pid"},
|
||||
}
|
||||
monkeypatch.setattr(agent_plugins, "liveness_for", lambda name: raw, raising=False)
|
||||
calls = []
|
||||
monkeypatch.setattr(redact, "register_vault_redaction_value", calls.append)
|
||||
caplog.set_level(logging.DEBUG)
|
||||
try:
|
||||
runtime.write_text(json.dumps({"http": "http://127.0.0.1:1111", "token": "first-secret", "pid": os.getpid()}))
|
||||
first = _live_endpoint("example-server")
|
||||
runtime.write_text(json.dumps({"http": "http://127.0.0.1:2222", "token": "second-secret", "pid": os.getpid()}))
|
||||
second = _live_endpoint("example-server")
|
||||
finally:
|
||||
declaration.unregister("example-server")
|
||||
assert first == ("http://127.0.0.1:1111/mcp", {"Authorization": "Bearer first-secret"})
|
||||
assert second == ("http://127.0.0.1:2222/mcp", {"Authorization": "Bearer second-secret"})
|
||||
assert calls == ["first-secret", "second-secret"]
|
||||
assert all(secret not in record.getMessage() for record in caplog.records for secret in calls)
|
||||
|
||||
|
||||
def test_runtime_file_without_token_connects_without_authorization(tmp_path, monkeypatch):
|
||||
import hermes_cli.agent_plugins as agent_plugins
|
||||
from agent import redact
|
||||
from tools.mcp_tool_transport import _live_endpoint
|
||||
|
||||
runtime = tmp_path / "server.json"
|
||||
declaration.register("example-server", _decl(tmp_path))
|
||||
monkeypatch.setattr(agent_plugins, "liveness_for", lambda name: {
|
||||
"kind": "server_json",
|
||||
"path": str(runtime),
|
||||
"fields": {"url": "http", "token": "token", "pid": "pid"},
|
||||
}, raising=False)
|
||||
calls = []
|
||||
monkeypatch.setattr(redact, "register_vault_redaction_value", calls.append)
|
||||
try:
|
||||
runtime.write_text(json.dumps({"http": "http://127.0.0.1:3333", "pid": os.getpid()}))
|
||||
result = _live_endpoint("example-server")
|
||||
finally:
|
||||
declaration.unregister("example-server")
|
||||
assert result is not None
|
||||
url, headers = result
|
||||
assert url == "http://127.0.0.1:3333/mcp"
|
||||
assert "Authorization" not in headers
|
||||
assert calls == []
|
||||
|
||||
|
||||
def test_missing_runtime_file_never_falls_back(tmp_path, monkeypatch):
|
||||
import hermes_cli.agent_plugins as agent_plugins
|
||||
from tools.mcp_tool_transport import LiveEndpointUnavailable, _live_endpoint
|
||||
|
||||
decl = _decl(tmp_path)
|
||||
declaration.register("example-server", decl)
|
||||
monkeypatch.setattr(agent_plugins, "liveness_for", lambda name: {
|
||||
"kind": "server_json",
|
||||
"path": str(tmp_path / "missing.json"),
|
||||
"fields": {"url": "http", "token": "token", "pid": "pid"},
|
||||
}, raising=False)
|
||||
try:
|
||||
with pytest.raises(LiveEndpointUnavailable):
|
||||
_live_endpoint("example-server")
|
||||
finally:
|
||||
declaration.unregister("example-server")
|
||||
|
||||
|
||||
def test_hydrated_error_shape_for_registered_declaration(tmp_path, monkeypatch):
|
||||
import hermes_cli.agent_plugins as agent_plugins
|
||||
from tools import mcp_tool, mcp_tool_discovery, mcp_tool_handlers
|
||||
|
||||
decl = _decl(tmp_path)
|
||||
declaration.register("example-server", decl)
|
||||
monkeypatch.setattr(agent_plugins, "liveness_for", lambda name: {"kind": "static"}, raising=False)
|
||||
monkeypatch.setattr(mcp_tool_discovery, "_get_connected_server_for_call", lambda name: None)
|
||||
monkeypatch.setattr(mcp_tool, "_bump_server_error", lambda name, **kwargs: None)
|
||||
try:
|
||||
server, error = mcp_tool_handlers._acquire_call_server("example-server", 0)
|
||||
finally:
|
||||
declaration.unregister("example-server")
|
||||
payload = json.loads(error)
|
||||
assert server is None
|
||||
assert payload["server"] == "example-server"
|
||||
assert payload["state"] == "app_not_running"
|
||||
assert payload["app"]["name"] == "Example App"
|
||||
assert payload["user_action"]
|
||||
assert payload["retry"] == "after_user_action"
|
||||
|
||||
|
||||
def test_connected_interactive_session_server_is_offerable_from_a_service_session(tmp_path, monkeypatch):
|
||||
import hermes_cli.agent_plugins as agent_plugins
|
||||
from hermes_platform.host import facts
|
||||
from tools import mcp_tool_handlers
|
||||
|
||||
declaration.register("example-server", _decl(tmp_path))
|
||||
monkeypatch.setattr(agent_plugins, "liveness_for", lambda name: {"kind": "interactive_session"}, raising=False)
|
||||
monkeypatch.setattr(facts, "interactive_session", lambda: False)
|
||||
try:
|
||||
assert mcp_tool_handlers._declared_app_offerable("example-server") is True
|
||||
finally:
|
||||
declaration.unregister("example-server")
|
||||
|
||||
|
||||
def test_undeclared_error_text_is_unchanged(monkeypatch):
|
||||
from tools import mcp_tool, mcp_tool_discovery, mcp_tool_handlers
|
||||
|
||||
monkeypatch.setattr(mcp_tool_discovery, "_get_connected_server_for_call", lambda name: None)
|
||||
monkeypatch.setattr(mcp_tool, "_bump_server_error", lambda name, **kwargs: None)
|
||||
_server, error = mcp_tool_handlers._acquire_call_server("plain-server", 0)
|
||||
assert json.loads(error)["error"] == "MCP server 'plain-server' is not connected"
|
||||
@@ -0,0 +1,74 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import sys
|
||||
|
||||
from hermes_platform import declaration
|
||||
from tools.registry import registry
|
||||
|
||||
|
||||
def _schema(name):
|
||||
return {"name": name, "description": "fixture", "parameters": {"type": "object", "properties": {}}}
|
||||
|
||||
|
||||
def _decl(tmp_path):
|
||||
missing = tmp_path / "missing-app"
|
||||
return declaration.parse_declaration(
|
||||
"Example App",
|
||||
{sys.platform: {"presence": "executable", "location": str(missing)}},
|
||||
{"app": True},
|
||||
where="test",
|
||||
)
|
||||
|
||||
|
||||
def test_hidden_declared_server_appears_in_listing_and_empty_summary(tmp_path, monkeypatch):
|
||||
import hermes_cli.agent_plugins as agent_plugins
|
||||
from tools.tool_search import ToolSearchConfig, dispatch_tool_search
|
||||
from tools.tool_search_catalog import build_catalog_listing_with_form
|
||||
|
||||
server = "example-hidden"
|
||||
tool_name = "mcp__example-hidden__inspect"
|
||||
declaration.register(server, _decl(tmp_path))
|
||||
monkeypatch.setattr(agent_plugins, "liveness_for", lambda name: {"kind": "static"}, raising=False)
|
||||
registry.register(
|
||||
name=tool_name,
|
||||
toolset=f"mcp-{server}",
|
||||
schema=_schema(tool_name),
|
||||
handler=lambda args: "{}",
|
||||
check_fn=lambda: False,
|
||||
)
|
||||
try:
|
||||
first = build_catalog_listing_with_form([], max_tokens=4000)
|
||||
second = build_catalog_listing_with_form([], max_tokens=4000)
|
||||
result = json.loads(dispatch_tool_search(
|
||||
{"queries": ["inspect"]},
|
||||
current_tool_defs=[],
|
||||
config=ToolSearchConfig(enabled="on", threshold_pct=50.0, search_default_limit=10, max_search_limit=50, defer_tools=frozenset({tool_name})),
|
||||
))
|
||||
finally:
|
||||
registry.deregister(tool_name)
|
||||
declaration.unregister(server)
|
||||
assert first == second
|
||||
assert first[1] == "full"
|
||||
assert "example-hidden (1 tools unavailable: Example App is not installed." in first[0]
|
||||
source = result["results"][0]["available_sources"][0]
|
||||
assert source["name"] == server and source["tool_count"] == 1
|
||||
assert "Example App is not installed." in source["unavailable"]
|
||||
|
||||
|
||||
def test_hidden_undeclared_server_is_absent(monkeypatch):
|
||||
from tools.tool_search_catalog import build_catalog_listing_with_form
|
||||
|
||||
tool_name = "mcp__plain-hidden__inspect"
|
||||
registry.register(
|
||||
name=tool_name,
|
||||
toolset="mcp-plain-hidden",
|
||||
schema=_schema(tool_name),
|
||||
handler=lambda args: "{}",
|
||||
check_fn=lambda: False,
|
||||
)
|
||||
try:
|
||||
listing, form = build_catalog_listing_with_form([], max_tokens=4000)
|
||||
finally:
|
||||
registry.deregister(tool_name)
|
||||
assert (listing, form) == (None, "none")
|
||||
@@ -0,0 +1,178 @@
|
||||
"""Application-backed MCP liveness parsing and user-facing status descriptions."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass, replace
|
||||
from typing import Any, Literal
|
||||
|
||||
from hermes_platform import declaration
|
||||
from hermes_platform.host import facts
|
||||
from hermes_platform.resolver.app import AppDef, AppResolver
|
||||
from hermes_platform.resolver.availability import Availability, availability
|
||||
from hermes_platform.resolver.base import Effort
|
||||
from hermes_platform.resolver.core import CheckState
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
LivenessKind = Literal["static", "server_json", "interactive_session"]
|
||||
LivenessState = Literal[
|
||||
"app_not_running",
|
||||
"endpoint_unavailable",
|
||||
"no_interactive_session",
|
||||
"version_too_old",
|
||||
"missing_app",
|
||||
]
|
||||
Retry = Literal["after_user_action", "never_here"]
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Liveness:
|
||||
kind: LivenessKind
|
||||
path: str = ""
|
||||
url_field: str = "http"
|
||||
token_field: str = "token"
|
||||
pid_field: str = "pid"
|
||||
|
||||
def app_definition(self, definition: AppDef) -> AppDef:
|
||||
if self.kind != "server_json":
|
||||
return definition
|
||||
return replace(
|
||||
definition,
|
||||
liveness_kind="server_json",
|
||||
liveness_path=self.path,
|
||||
liveness_pid_key=self.pid_field,
|
||||
liveness_url_key=self.url_field,
|
||||
liveness_token_key=self.token_field,
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Status:
|
||||
state: LivenessState
|
||||
availability: Availability
|
||||
liveness: Liveness
|
||||
user_action: str
|
||||
retry: Retry
|
||||
|
||||
|
||||
def parse_liveness(raw: Any) -> Liveness:
|
||||
"""Parse one portable-plugin liveness declaration."""
|
||||
if not isinstance(raw, dict):
|
||||
raise ValueError("liveness must be an object")
|
||||
kind = raw.get("kind")
|
||||
if kind == "static":
|
||||
if set(raw) != {"kind"}:
|
||||
raise ValueError("static liveness only accepts 'kind'")
|
||||
return Liveness("static")
|
||||
if kind == "interactive_session":
|
||||
if set(raw) != {"kind"}:
|
||||
raise ValueError("interactive_session liveness only accepts 'kind'")
|
||||
return Liveness("interactive_session")
|
||||
if kind != "server_json":
|
||||
raise ValueError(f"unknown liveness kind: {kind!r}")
|
||||
if set(raw) - {"kind", "path", "fields"}:
|
||||
raise ValueError("server_json liveness has unknown fields")
|
||||
path = raw.get("path")
|
||||
fields = {"url": "http", "token": "token", "pid": "pid", **(raw.get("fields") or {})}
|
||||
if not isinstance(path, str) or not path.strip():
|
||||
raise ValueError("server_json liveness requires a non-empty path")
|
||||
if set(fields) != {"url", "token", "pid"}:
|
||||
raise ValueError("server_json liveness fields may only override url, token, and pid")
|
||||
if any(not isinstance(value, str) or not value.strip() for value in fields.values()):
|
||||
raise ValueError("server_json liveness field names must be non-empty strings")
|
||||
return Liveness(
|
||||
"server_json",
|
||||
path=path.strip(),
|
||||
url_field=fields["url"].strip(),
|
||||
token_field=fields["token"].strip(),
|
||||
pid_field=fields["pid"].strip(),
|
||||
)
|
||||
|
||||
|
||||
def liveness_for(server_name: str) -> Liveness:
|
||||
"""Return a server's registered liveness declaration, defaulting to static."""
|
||||
try:
|
||||
from hermes_cli.agent_plugins import liveness_for as registered_liveness
|
||||
except ImportError:
|
||||
return Liveness("static")
|
||||
raw = registered_liveness(server_name)
|
||||
if raw is None:
|
||||
return Liveness("static")
|
||||
try:
|
||||
return parse_liveness(raw)
|
||||
except ValueError as exc:
|
||||
logger.warning("MCP server '%s' has an invalid liveness declaration (%s); treating it as static", server_name, exc)
|
||||
return Liveness("static")
|
||||
|
||||
|
||||
def _action(state: LivenessState, app_name: str) -> tuple[str, Retry]:
|
||||
actions: dict[LivenessState, tuple[str, Retry]] = {
|
||||
"app_not_running": (f"Start {app_name}, then try again.", "after_user_action"),
|
||||
"endpoint_unavailable": (f"Open {app_name} and enable its local connection, then try again.", "after_user_action"),
|
||||
"no_interactive_session": (f"Open an interactive desktop session and start {app_name}, then try again.", "never_here"),
|
||||
"version_too_old": (f"Update {app_name}, then try again.", "after_user_action"),
|
||||
"missing_app": (f"Install {app_name}, then try again.", "after_user_action"),
|
||||
}
|
||||
return actions[state]
|
||||
|
||||
|
||||
def status(server_name: str) -> Status | None:
|
||||
"""Return the current unavailable state for a registered declaration."""
|
||||
decl = declaration.lookup(server_name)
|
||||
if decl is None:
|
||||
return None
|
||||
available = availability(decl)
|
||||
live = liveness_for(server_name)
|
||||
if available.state in {"missing_app", "unsupported_os"}:
|
||||
state: LivenessState = "missing_app"
|
||||
elif available.state == "version_too_old":
|
||||
state = "version_too_old"
|
||||
elif live.kind == "interactive_session" and not facts.interactive_session():
|
||||
state = "no_interactive_session"
|
||||
elif live.kind == "server_json":
|
||||
definition = decl.app_for(facts.os_family())
|
||||
if definition is None:
|
||||
state = "missing_app"
|
||||
else:
|
||||
probe = AppResolver(live.app_definition(definition)).probe(
|
||||
AppResolver(live.app_definition(definition)).locate(), effort=Effort.LOCAL
|
||||
)
|
||||
if probe.running.value is not True:
|
||||
state = "app_not_running"
|
||||
elif probe.endpoint.state is not CheckState.PRESENT:
|
||||
state = "endpoint_unavailable"
|
||||
else:
|
||||
state = "app_not_running"
|
||||
else:
|
||||
state = "app_not_running"
|
||||
action, retry = _action(state, decl.name)
|
||||
return Status(state, available, live, action, retry)
|
||||
|
||||
|
||||
def describe(decl: declaration.Declaration, available: Availability, liveness_state: LivenessState) -> str:
|
||||
"""Compose one unavailable-state sentence with one user action."""
|
||||
app_name = decl.name
|
||||
action, _retry = _action(liveness_state, app_name)
|
||||
if liveness_state == "missing_app":
|
||||
reason = f"{app_name} is not installed."
|
||||
elif liveness_state == "version_too_old":
|
||||
found = f" version {available.version}" if available.version else ""
|
||||
minimum = f"; version {available.min_version} or newer is required" if available.min_version else ""
|
||||
reason = f"{app_name}{found} is too old{minimum}."
|
||||
elif liveness_state == "no_interactive_session":
|
||||
reason = f"{app_name} needs an interactive desktop session."
|
||||
elif liveness_state == "endpoint_unavailable":
|
||||
reason = f"{app_name}'s local endpoint is unavailable."
|
||||
else:
|
||||
reason = f"{app_name} is not running."
|
||||
return f"{reason} {action}"
|
||||
|
||||
|
||||
def unavailable_details(server_name: str) -> tuple[declaration.Declaration, Status, str] | None:
|
||||
"""Return declaration, structured state, and its composed sentence."""
|
||||
decl = declaration.lookup(server_name)
|
||||
current = status(server_name)
|
||||
if decl is None or current is None:
|
||||
return None
|
||||
return decl, current, describe(decl, current.availability, current.state)
|
||||
@@ -111,6 +111,22 @@ def _acquire_call_server(server_name: str, tool_timeout: float):
|
||||
server task to rebuild (probing a dead transport would re-arm the breaker forever)."""
|
||||
from tools import mcp_tool_discovery as _discovery # lazy: discovery -> registration -> handlers cycle
|
||||
not_connected = tool_error(f"MCP server '{server_name}' is not connected")
|
||||
from tools.mcp_liveness import unavailable_details
|
||||
details = unavailable_details(server_name)
|
||||
if details is not None:
|
||||
decl, current, sentence = details
|
||||
not_connected = tool_error(
|
||||
sentence,
|
||||
server=server_name,
|
||||
state=current.state,
|
||||
app={
|
||||
"name": decl.name,
|
||||
"version": current.availability.version,
|
||||
"path": current.availability.path,
|
||||
},
|
||||
user_action=current.user_action,
|
||||
retry=current.retry,
|
||||
)
|
||||
server = _discovery._get_connected_server_for_call(server_name)
|
||||
wait = min(5.0, float(tool_timeout or 5.0))
|
||||
if server and (server.session or _loop._wait_for_server_session_ready(server, timeout=wait)):
|
||||
@@ -695,11 +711,12 @@ def _make_check_fn(server_name: str):
|
||||
|
||||
|
||||
def _declared_app_offerable(server_name: str) -> bool:
|
||||
"""True unless a declaration registered for this server requires an application this host lacks."""
|
||||
"""True unless the registered declaration is unavailable on this host. Called only for a
|
||||
connected server, so a reachable loopback port outranks the interactive-session rule."""
|
||||
from hermes_platform import declaration
|
||||
from hermes_platform.resolver.availability import availability
|
||||
|
||||
decl = declaration.lookup(server_name)
|
||||
if decl is None or not decl.requires_app:
|
||||
if decl is None:
|
||||
return True
|
||||
return availability(decl).offerable
|
||||
return bool(availability(decl).offerable)
|
||||
|
||||
@@ -284,7 +284,9 @@ class MCPServerRunMixin:
|
||||
# Content-type preflight (Streamable HTTP only; SSE serves text/event-stream): a
|
||||
# web-app root returns HTML and would hang the SDK for connect_timeout. Skipped once
|
||||
# _ready was ever set and for OAuth servers (a token-less probe sees HTML/401).
|
||||
from tools.mcp_liveness import liveness_for
|
||||
if (config.get("transport") != "sse" and not config.get("skip_preflight")
|
||||
and liveness_for(self.name).kind != "server_json"
|
||||
and not self._ready.is_set() and self._auth_type != "oauth"):
|
||||
await self._preflight_content_type(
|
||||
config["url"], headers=dict(config.get("headers") or {}),
|
||||
|
||||
@@ -86,6 +86,31 @@ def _pgroup_alive(pgid: Optional[int]) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
class LiveEndpointUnavailable(ConnectionError):
|
||||
"""A declared runtime file did not provide a usable live endpoint."""
|
||||
|
||||
|
||||
def _live_endpoint(server_name: str) -> Optional[tuple[str, dict]]:
|
||||
from agent.redact import register_vault_redaction_value
|
||||
from hermes_platform import declaration
|
||||
from hermes_platform.host import facts
|
||||
from hermes_platform.resolver.app import AppResolver
|
||||
from tools.mcp_liveness import liveness_for
|
||||
|
||||
live = liveness_for(server_name)
|
||||
if live.kind != "server_json":
|
||||
return None
|
||||
decl = declaration.lookup(server_name)
|
||||
definition = decl.app_for(facts.os_family()) if decl is not None else None
|
||||
endpoint = AppResolver(live.app_definition(definition)).endpoint() if definition is not None else None
|
||||
if endpoint is None:
|
||||
raise LiveEndpointUnavailable(f"MCP server '{server_name}' has no usable live endpoint")
|
||||
if endpoint.token:
|
||||
register_vault_redaction_value(endpoint.token)
|
||||
headers = {"Authorization": f"Bearer {endpoint.token}"} if endpoint.token else {}
|
||||
return endpoint.url, headers
|
||||
|
||||
|
||||
class MCPServerTransportMixin:
|
||||
"""Methods of :class:`tools.mcp_tool.MCPServerTask` (mixed in; relies on its attributes)."""
|
||||
|
||||
@@ -504,9 +529,13 @@ class MCPServerTransportMixin:
|
||||
"mcp.client.streamable_http is not available. "
|
||||
"Upgrade the mcp package to get HTTP support.")
|
||||
url = config["url"]
|
||||
headers = dict(config.get("headers") or {})
|
||||
live = _live_endpoint(self.name)
|
||||
if live is not None:
|
||||
url, live_headers = live
|
||||
headers.update(live_headers)
|
||||
logger.debug("MCP server '%s': connecting to %s", self.name, url)
|
||||
self._http_rejection = {} # last 4xx/5xx the owned client saw this attempt (recorder hook)
|
||||
headers = dict(config.get("headers") or {})
|
||||
# Agent Plugins v1 strict_redirect_headers: configured headers MUST NOT follow a cross-origin
|
||||
# redirect — capture their names BEFORE client-generated headers are merged in.
|
||||
configured_header_names = {key.lower() for key in headers}
|
||||
|
||||
@@ -409,10 +409,12 @@ def _shared_tool_record(entry: CatalogEntry) -> Dict[str, Any]:
|
||||
|
||||
|
||||
def _available_source_summary(catalog: List[CatalogEntry]) -> List[Dict[str, Any]]:
|
||||
"""Deterministic ``[{name, tool_count}]`` of connected sources (attached to empty query
|
||||
groups so a lexical miss is not read as a missing capability)."""
|
||||
"""Deterministic summaries of connected and declared unavailable sources."""
|
||||
from tools.tool_search_catalog import hidden_declared_sources
|
||||
|
||||
counts = Counter(_listing_group_label(entry.source_name) for entry in catalog)
|
||||
return [{"name": name, "tool_count": counts[name]} for name in sorted(counts)]
|
||||
rows = [{"name": name, "tool_count": counts[name]} for name in sorted(counts)]
|
||||
return sorted(rows + hidden_declared_sources(), key=lambda row: row["name"])
|
||||
|
||||
|
||||
def _string_list_arg(args: Dict[str, Any], key: str, *, dedupe: bool, max_items: int,
|
||||
@@ -454,7 +456,7 @@ def dispatch_tool_search(args: Dict[str, Any], *, current_tool_defs: List[Dict[s
|
||||
queries, connector_search=connector_search)
|
||||
results: List[Dict[str, Any]] = []
|
||||
tools_map: Dict[str, Dict[str, Any]] = {}
|
||||
available_sources = _available_source_summary(catalog) if catalog else []
|
||||
available_sources = _available_source_summary(catalog)
|
||||
for position, query in enumerate(queries):
|
||||
corpus = catalog + remote_entries[position]
|
||||
hits = search_catalog(corpus, query, limit=limit)
|
||||
@@ -462,7 +464,7 @@ def dispatch_tool_search(args: Dict[str, Any], *, current_tool_defs: List[Dict[s
|
||||
tools_map.setdefault(h.name, _shared_tool_record(h))
|
||||
matches = [h.name for h in hits]
|
||||
group: Dict[str, Any] = {"query": query, "matches": matches}
|
||||
if not matches and catalog:
|
||||
if not matches and available_sources:
|
||||
group["available_sources"] = available_sources
|
||||
group["hint"] = (
|
||||
"This query returned no lexical matches, but the sources above "
|
||||
|
||||
@@ -233,6 +233,31 @@ def _listing_group_label(source_name: str) -> str:
|
||||
return label[4:] if label.startswith("mcp-") else label
|
||||
|
||||
|
||||
def hidden_declared_sources() -> List[Dict[str, Any]]:
|
||||
"""Return deterministic summaries for declared MCP servers hidden by their check."""
|
||||
from hermes_platform import declaration
|
||||
from tools.mcp_liveness import unavailable_details
|
||||
from tools.registry import registry
|
||||
|
||||
grouped: Dict[str, List[Any]] = {}
|
||||
for entry in registry.get_all_entries():
|
||||
if entry.toolset.startswith("mcp-"):
|
||||
grouped.setdefault(entry.toolset[4:], []).append(entry)
|
||||
rows: List[Dict[str, Any]] = []
|
||||
for server_name in sorted(grouped):
|
||||
if declaration.lookup(server_name) is None:
|
||||
continue
|
||||
entries = grouped[server_name]
|
||||
if any(entry.check_fn is None or bool(entry.check_fn()) for entry in entries):
|
||||
continue
|
||||
details = unavailable_details(server_name)
|
||||
if details is None:
|
||||
continue
|
||||
_decl, _current, sentence = details
|
||||
rows.append({"name": server_name, "tool_count": len(entries), "unavailable": sentence})
|
||||
return rows
|
||||
|
||||
|
||||
def build_catalog_listing_with_form(
|
||||
deferrable: List[Dict[str, Any]], *, max_tokens: int = 4000) -> Tuple[Optional[str], str]:
|
||||
"""Render the deferred-catalog manifest: ``- name: short desc`` lines grouped per source.
|
||||
@@ -249,7 +274,8 @@ def build_catalog_listing_with_form(
|
||||
# _classify_source gives ("other", "") when unregistered; the label of "" is "other".
|
||||
label = _listing_group_label(_classify_source(name)[1])
|
||||
groups.setdefault(label, []).append((name, _short_desc(fn.get("description", ""))))
|
||||
if not groups:
|
||||
unavailable = hidden_declared_sources()
|
||||
if not groups and not unavailable:
|
||||
return None, "none"
|
||||
|
||||
def render_group(label: str, mode: str) -> str:
|
||||
@@ -269,7 +295,20 @@ def build_catalog_listing_with_form(
|
||||
f"`{TOOL_DESCRIBE_NAME}`, invoke via `{TOOL_CALL_NAME}`):")
|
||||
|
||||
def assemble_if_fits(modes: Dict[str, str]) -> Optional[str]:
|
||||
text = "\n".join([header] + [render_group(lbl, modes[lbl]) for lbl in sorted(groups)])
|
||||
available_blocks = {label: render_group(label, modes[label]) for label in groups}
|
||||
unavailable_blocks = {
|
||||
row["name"]: (
|
||||
f"{row['name']} ({row['tool_count']} tools unavailable: {row['unavailable']})"
|
||||
if row.get("tool_count") is not None
|
||||
else f"{row['name']} (tools unavailable: {row['unavailable']})"
|
||||
)
|
||||
for row in unavailable
|
||||
}
|
||||
blocks = [
|
||||
available_blocks[label] if label in available_blocks else unavailable_blocks[label]
|
||||
for label in sorted(available_blocks | unavailable_blocks)
|
||||
]
|
||||
text = "\n".join([header] + blocks)
|
||||
return text if math.ceil(len(text) / CHARS_PER_TOKEN) <= max_tokens else None
|
||||
|
||||
for mode in ("full", "names"): # 1. everything full; 2. everything names-only
|
||||
|
||||
Reference in New Issue
Block a user