300 lines
12 KiB
Python
300 lines
12 KiB
Python
#!/usr/bin/env python3
|
|
"""Stage portable MCP runtimes in persistent MCPHub Appdata safely.
|
|
|
|
The helper deliberately separates staging from activation. Missing credentials
|
|
can never result in a published server, even if an agent submits a manifest
|
|
that requests activation.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import hashlib
|
|
import json
|
|
import os
|
|
import pathlib
|
|
import re
|
|
import shutil
|
|
import tempfile
|
|
import urllib.error
|
|
import urllib.parse
|
|
import urllib.request
|
|
|
|
|
|
ID_RE = re.compile(r"^[a-z0-9][a-z0-9-]{0,62}$")
|
|
|
|
|
|
def load_json(path: pathlib.Path) -> dict:
|
|
value = json.loads(path.read_text(encoding="utf-8"))
|
|
if not isinstance(value, dict):
|
|
raise SystemExit(f"Expected JSON object: {path}")
|
|
return value
|
|
|
|
|
|
def atomic_json(path: pathlib.Path, value: dict) -> None:
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
fd, temporary = tempfile.mkstemp(prefix=f".{path.name}.", dir=path.parent)
|
|
try:
|
|
with os.fdopen(fd, "w", encoding="utf-8") as handle:
|
|
json.dump(value, handle, indent=2, ensure_ascii=False)
|
|
handle.write("\n")
|
|
os.chmod(temporary, 0o600)
|
|
os.replace(temporary, path)
|
|
finally:
|
|
if os.path.exists(temporary):
|
|
os.unlink(temporary)
|
|
|
|
|
|
def env_keys(path: pathlib.Path) -> set[str]:
|
|
keys: set[str] = set()
|
|
if not path.is_file():
|
|
return keys
|
|
for raw in path.read_text(encoding="utf-8", errors="replace").splitlines():
|
|
line = raw.strip()
|
|
if not line or line.startswith("#") or "=" not in line:
|
|
continue
|
|
key, value = line.split("=", 1)
|
|
if value.strip().strip("\"'"):
|
|
keys.add(key.removeprefix("export ").strip())
|
|
return keys
|
|
|
|
|
|
def credential_state(server: dict, secrets_dir: pathlib.Path) -> tuple[bool, str]:
|
|
deployment = server.get("deployment") or {}
|
|
required = [str(item) for item in deployment.get("required_env", [])]
|
|
required_files = [str(item) for item in deployment.get("required_files", [])]
|
|
secret_name = str((server.get("hub") or {}).get("secret_file") or "")
|
|
if not required and not secret_name:
|
|
return True, "not-required"
|
|
if not secret_name:
|
|
return False, "secret-file-not-declared"
|
|
source = secrets_dir / secret_name
|
|
present = env_keys(source)
|
|
missing = [key for key in required if key not in present]
|
|
if not source.is_file():
|
|
return False, f"missing:{source}"
|
|
if missing:
|
|
return False, "missing-keys:" + ",".join(missing)
|
|
missing_files = [name for name in required_files if not (secrets_dir / name).is_file()]
|
|
if missing_files:
|
|
return False, "missing-files:" + ",".join(missing_files)
|
|
return True, "ready"
|
|
|
|
|
|
def registry(path: pathlib.Path) -> dict:
|
|
document = load_json(path)
|
|
if document.get("version") != 1 or not isinstance(document.get("servers"), list):
|
|
raise SystemExit("Unsupported MCP registry schema")
|
|
return document
|
|
|
|
|
|
def find_server(document: dict, server_id: str) -> dict | None:
|
|
return next((item for item in document["servers"] if item.get("id") == server_id), None)
|
|
|
|
|
|
def expand_runtime(value: object, values: dict[str, str]) -> object:
|
|
if isinstance(value, str):
|
|
def replace(match: re.Match[str]) -> str:
|
|
key = match.group(1)
|
|
resolved = values.get(key, "").strip()
|
|
if not resolved:
|
|
raise SystemExit(f"Missing runtime value: {key}")
|
|
return resolved
|
|
return re.sub(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}", replace, value)
|
|
if isinstance(value, list):
|
|
return [expand_runtime(item, values) for item in value]
|
|
if isinstance(value, dict):
|
|
return {key: expand_runtime(item, values) for key, item in value.items()}
|
|
return value
|
|
|
|
|
|
def api_request(args: argparse.Namespace, method: str, path: str,
|
|
body: dict | None = None, allow_not_found: bool = False) -> dict | None:
|
|
if getattr(args, "skip_api", False):
|
|
return None
|
|
token_file = getattr(args, "token_file", None) or args.appdata / "client-token"
|
|
token = pathlib.Path(token_file).read_text(encoding="utf-8").strip()
|
|
if not token:
|
|
raise SystemExit("MCPHub API token is empty")
|
|
payload = None if body is None else json.dumps(body).encode("utf-8")
|
|
request = urllib.request.Request(
|
|
str(getattr(args, "api_url", "http://127.0.0.1:3000/api")).rstrip("/") + path,
|
|
data=payload,
|
|
method=method,
|
|
headers={
|
|
"Authorization": f"Bearer {token}",
|
|
"Accept": "application/json",
|
|
"Content-Type": "application/json",
|
|
},
|
|
)
|
|
try:
|
|
with urllib.request.urlopen(request, timeout=30) as response:
|
|
raw = response.read(1_000_000)
|
|
except urllib.error.HTTPError as exc:
|
|
if allow_not_found and exc.code == 404:
|
|
return None
|
|
detail = exc.read(500).decode("utf-8", errors="replace").replace("\n", " ")
|
|
raise SystemExit(f"MCPHub API HTTP {exc.code}: {detail[:300]}") from exc
|
|
except OSError as exc:
|
|
raise SystemExit(f"MCPHub API unavailable: {str(exc)[:300]}") from exc
|
|
return json.loads(raw) if raw else {}
|
|
|
|
|
|
def sync_runtime(args: argparse.Namespace, server: dict, enabled: bool,
|
|
credentials_ready: bool) -> None:
|
|
if getattr(args, "skip_api", False):
|
|
return
|
|
name = str(server.get("hermes_id") or server["id"])
|
|
config = json.loads(json.dumps(server["hub"]))
|
|
secret_name = str(config.pop("secret_file", ""))
|
|
# env_keys intentionally returns names only. Re-read values only for the
|
|
# runtime rendering path, never print or return them.
|
|
if credentials_ready and secret_name:
|
|
values: dict[str, str] = {}
|
|
for raw in (args.secrets / secret_name).read_text(encoding="utf-8", errors="replace").splitlines():
|
|
line = raw.strip()
|
|
if not line or line.startswith("#") or "=" not in line:
|
|
continue
|
|
key, value = line.split("=", 1)
|
|
values[key.removeprefix("export ").strip()] = value.strip().strip("\"'")
|
|
config = expand_runtime(json.loads(json.dumps(server["hub"])), values)
|
|
config.pop("secret_file", None)
|
|
config["enabled"] = bool(enabled)
|
|
current = api_request(args, "GET", f"/servers/{urllib.parse.quote(name, safe='')}", allow_not_found=True)
|
|
if current is None:
|
|
api_request(args, "POST", "/servers", {"name": name, "config": config})
|
|
else:
|
|
api_request(args, "PUT", f"/servers/{urllib.parse.quote(name, safe='')}", {"config": config})
|
|
|
|
|
|
def validate_server(server: dict) -> str:
|
|
server_id = str(server.get("id") or "")
|
|
if not ID_RE.fullmatch(server_id):
|
|
raise SystemExit("Invalid server id")
|
|
for key in ("name", "description", "url", "hub"):
|
|
if not server.get(key):
|
|
raise SystemExit(f"Server field is required: {key}")
|
|
if not isinstance(server["hub"], dict) or not server["hub"].get("type"):
|
|
raise SystemExit("hub.type is required")
|
|
return server_id
|
|
|
|
|
|
def stage(args: argparse.Namespace) -> None:
|
|
manifest = load_json(args.manifest)
|
|
server = manifest.get("server")
|
|
if not isinstance(server, dict):
|
|
raise SystemExit("manifest.server must be an object")
|
|
server = json.loads(json.dumps(server))
|
|
server_id = validate_server(server)
|
|
extension_dir = args.appdata / "extensions" / server_id
|
|
work_root = (args.appdata / "work").resolve()
|
|
extension_dir.parent.mkdir(parents=True, exist_ok=True)
|
|
temporary = pathlib.Path(tempfile.mkdtemp(prefix=f".{server_id}.", dir=extension_dir.parent))
|
|
try:
|
|
for artifact in manifest.get("artifacts", []):
|
|
source = pathlib.Path(str(artifact["source"])).resolve()
|
|
if work_root not in source.parents:
|
|
raise SystemExit(f"Artifact must be under {work_root}")
|
|
relative = pathlib.PurePosixPath(str(artifact["path"]))
|
|
if relative.is_absolute() or ".." in relative.parts:
|
|
raise SystemExit("Invalid artifact destination")
|
|
expected = str(artifact["sha256"]).lower()
|
|
actual = hashlib.sha256(source.read_bytes()).hexdigest()
|
|
if actual != expected:
|
|
raise SystemExit(f"Checksum mismatch for {relative}")
|
|
target = temporary / relative
|
|
target.parent.mkdir(parents=True, exist_ok=True)
|
|
shutil.copyfile(source, target)
|
|
target.chmod(int(str(artifact.get("mode", "0644")), 8))
|
|
backup = extension_dir.with_name(extension_dir.name + ".previous")
|
|
if backup.exists():
|
|
shutil.rmtree(backup)
|
|
if extension_dir.exists():
|
|
extension_dir.rename(backup)
|
|
temporary.rename(extension_dir)
|
|
except BaseException:
|
|
shutil.rmtree(temporary, ignore_errors=True)
|
|
raise
|
|
|
|
desired = list(server.get("clients", []))
|
|
deployment = server.setdefault("deployment", {})
|
|
deployment["desired_clients"] = desired
|
|
deployment["managed_enabled"] = True
|
|
ready, reason = credential_state(server, args.secrets)
|
|
server["hub"]["enabled"] = False
|
|
server["clients"] = []
|
|
document = registry(args.registry)
|
|
document["servers"] = [item for item in document["servers"] if item.get("id") != server_id]
|
|
document["servers"].append(server)
|
|
atomic_json(args.registry, document)
|
|
sync_runtime(args, server, False, ready)
|
|
print(json.dumps({
|
|
"status": "staged", "id": server_id, "enabled": False,
|
|
"credentials_ready": ready, "credential_state": reason,
|
|
"extension": str(extension_dir),
|
|
}))
|
|
|
|
|
|
def set_enabled(args: argparse.Namespace, enabled: bool) -> None:
|
|
document = registry(args.registry)
|
|
server = find_server(document, args.id)
|
|
if server is None:
|
|
raise SystemExit(f"Unknown server: {args.id}")
|
|
ready, reason = credential_state(server, args.secrets)
|
|
if enabled and not ready:
|
|
raise SystemExit(f"Activation refused: {reason}")
|
|
sync_runtime(args, server, enabled, ready)
|
|
server["hub"]["enabled"] = enabled
|
|
desired = list((server.get("deployment") or {}).get("desired_clients", []))
|
|
server["clients"] = desired if enabled else []
|
|
atomic_json(args.registry, document)
|
|
print(json.dumps({"status": "enabled" if enabled else "disabled", "id": args.id}))
|
|
|
|
|
|
def status(args: argparse.Namespace) -> None:
|
|
document = registry(args.registry)
|
|
server = find_server(document, args.id)
|
|
if server is None:
|
|
print(json.dumps({"id": args.id, "registered": False}))
|
|
return
|
|
ready, reason = credential_state(server, args.secrets)
|
|
print(json.dumps({
|
|
"id": args.id,
|
|
"registered": True,
|
|
"enabled": bool((server.get("hub") or {}).get("enabled")),
|
|
"clients": server.get("clients", []),
|
|
"credentials_ready": ready,
|
|
"credential_state": reason,
|
|
"extension_exists": (args.appdata / "extensions" / args.id).is_dir(),
|
|
}))
|
|
|
|
|
|
def main() -> None:
|
|
parser = argparse.ArgumentParser()
|
|
parser.add_argument("--appdata", type=pathlib.Path, default=pathlib.Path("/app/data"))
|
|
parser.add_argument("--registry", type=pathlib.Path, default=pathlib.Path("/app/data/config/mcp-registry.json"))
|
|
parser.add_argument("--secrets", type=pathlib.Path, default=pathlib.Path("/run/secrets/mcphub"))
|
|
parser.add_argument("--api-url", default="http://127.0.0.1:3000/api")
|
|
parser.add_argument("--token-file", type=pathlib.Path, default=pathlib.Path("/app/data/client-token"))
|
|
parser.add_argument("--skip-api", action="store_true", help=argparse.SUPPRESS)
|
|
sub = parser.add_subparsers(dest="command", required=True)
|
|
stage_cmd = sub.add_parser("stage")
|
|
stage_cmd.add_argument("--manifest", type=pathlib.Path, required=True)
|
|
for name in ("status", "activate", "disable"):
|
|
command = sub.add_parser(name)
|
|
command.add_argument("id")
|
|
args = parser.parse_args()
|
|
args.appdata.mkdir(parents=True, exist_ok=True)
|
|
if args.command == "stage":
|
|
stage(args)
|
|
elif args.command == "status":
|
|
status(args)
|
|
elif args.command == "activate":
|
|
set_enabled(args, True)
|
|
else:
|
|
set_enabled(args, False)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|