Add staged MCP release imports
This commit is contained in:
@@ -28,7 +28,7 @@ from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
|
||||
VERSION = "2.2.0"
|
||||
VERSION = "2.3.0"
|
||||
STACK = Path(os.environ.get("ATHENA_OPERATOR_STACK", "/opt/mike-ai/stack")).resolve()
|
||||
REPOSITORY = Path(os.environ.get("ATHENA_OPERATOR_REPOSITORY", "/data/mike-ai-operator/repository")).resolve()
|
||||
STATE = Path(os.environ.get("ATHENA_OPERATOR_STATE", "/data/mike-ai-operator/state")).resolve()
|
||||
@@ -39,6 +39,14 @@ MAX_FILE_BYTES = 1_000_000
|
||||
MAX_FILES = 24
|
||||
MAX_OUTPUT = 30_000
|
||||
LOCK = threading.RLock()
|
||||
STAGING_ROOTS = tuple(
|
||||
Path(value).resolve()
|
||||
for value in os.environ.get(
|
||||
"ATHENA_OPERATOR_STAGING_ROOTS",
|
||||
"/data/mike-ai-operator/staging:/data/hermes/deemix-mcp-build",
|
||||
).split(":")
|
||||
if value
|
||||
)
|
||||
|
||||
SAFE_PATH = re.compile(r"^[A-Za-z0-9_.+/-]{1,240}$")
|
||||
SAFE_NAME = re.compile(r"^[A-Za-z0-9_.-]{1,100}$")
|
||||
@@ -296,8 +304,59 @@ def search_source(arguments: dict[str, Any]) -> dict[str, Any]:
|
||||
query = str(arguments.get("query", ""))
|
||||
if not query or len(query) > 200 or any(x in query for x in ("\x00", "\n", "\r")):
|
||||
raise ValueError("invalid query")
|
||||
result = run(["rg", "-n", "--hidden", "--glob", "!.git/**", "--", query, "."], cwd=REPOSITORY, timeout=20)
|
||||
return {"query": query, "matches": result["output"], "exit_code": result["exit_code"]}
|
||||
if shutil.which("rg"):
|
||||
result = run(["rg", "-n", "--hidden", "--glob", "!.git/**", "--", query, "."], cwd=REPOSITORY, timeout=20)
|
||||
return {"query": query, "matches": result["output"], "exit_code": result["exit_code"], "engine": "rg"}
|
||||
try:
|
||||
pattern = re.compile(query)
|
||||
except re.error as exc:
|
||||
raise ValueError(f"invalid search expression: {exc}") from exc
|
||||
matches: list[str] = []
|
||||
for path in sorted(REPOSITORY.rglob("*")):
|
||||
if not path.is_file() or ".git" in path.parts or path.stat().st_size > MAX_FILE_BYTES:
|
||||
continue
|
||||
try:
|
||||
lines = path.read_text(encoding="utf-8", errors="strict").splitlines()
|
||||
except (UnicodeDecodeError, OSError):
|
||||
continue
|
||||
relative = path.relative_to(REPOSITORY)
|
||||
for number, line in enumerate(lines, 1):
|
||||
if pattern.search(line):
|
||||
matches.append(f"{relative}:{number}:{line}")
|
||||
if len(matches) >= 200:
|
||||
return {"query": query, "matches": "\n".join(matches), "exit_code": 0, "engine": "python", "truncated": True}
|
||||
return {"query": query, "matches": "\n".join(matches), "exit_code": 0 if matches else 1, "engine": "python", "truncated": False}
|
||||
|
||||
|
||||
def staged_file(item: dict[str, Any]) -> dict[str, Any]:
|
||||
source = Path(str(item.get("source", "")))
|
||||
if not source.is_absolute() or not source.is_file():
|
||||
raise FileNotFoundError("staged source file not found")
|
||||
resolved = source.resolve()
|
||||
if not any(resolved.is_relative_to(root) for root in STAGING_ROOTS):
|
||||
raise PermissionError("staged source is outside approved staging roots")
|
||||
raw = resolved.read_bytes()
|
||||
if len(raw) > MAX_FILE_BYTES:
|
||||
raise ValueError("staged source exceeds limit")
|
||||
expected = str(item.get("expected_source_sha256", ""))
|
||||
actual = sha(raw)
|
||||
if not expected or expected != actual:
|
||||
raise RuntimeError("staged source checksum is missing or does not match")
|
||||
try:
|
||||
content = raw.decode("utf-8", errors="strict")
|
||||
except UnicodeDecodeError as exc:
|
||||
raise ValueError("staged source must be UTF-8 text") from exc
|
||||
relative = safe_relative(str(item.get("path", "")))
|
||||
target = source_file(REPOSITORY, relative)
|
||||
before = target.read_text(encoding="utf-8", errors="strict") if target.is_file() else ""
|
||||
before_sha = sha(before.encode())
|
||||
expected_target = str(item.get("expected_target_sha256", ""))
|
||||
if expected_target and expected_target != before_sha:
|
||||
raise RuntimeError(f"source drift for {relative}")
|
||||
return {
|
||||
"path": str(relative), "content": content, "before_sha256": before_sha,
|
||||
"after_sha256": actual, "staged_source": str(resolved),
|
||||
}
|
||||
|
||||
|
||||
def terminal(arguments: dict[str, Any]) -> dict[str, Any]:
|
||||
@@ -362,7 +421,19 @@ def normalise_operation(operation: str, payload: dict[str, Any]) -> tuple[dict[s
|
||||
files, preview = normalise_patches(payload.get("files"))
|
||||
return {"files": files}, preview
|
||||
if operation == "mcp_release":
|
||||
files, patch_preview = normalise_patches(payload.get("files"))
|
||||
patch_items = payload.get("files") or []
|
||||
import_items = payload.get("imports") or []
|
||||
if not isinstance(patch_items, list) or not isinstance(import_items, list):
|
||||
raise ValueError("files and imports must be lists")
|
||||
if not patch_items and not import_items:
|
||||
raise ValueError("mcp_release requires at least one patch or staged import")
|
||||
files, patch_preview = normalise_patches(patch_items) if patch_items else ([], "")
|
||||
imports = [staged_file(item) for item in import_items]
|
||||
files.extend(imports)
|
||||
import_preview = "\n".join(
|
||||
f"IMPORT {item['staged_source']} -> {item['path']} sha256={item['after_sha256']}"
|
||||
for item in imports
|
||||
)
|
||||
checks = payload.get("checks") or ["operator-tests", "compose-mcp"]
|
||||
if not isinstance(checks, list) or not checks or any(name not in ALLOWED_CHECKS for name in checks):
|
||||
raise ValueError("unknown release check suite")
|
||||
@@ -386,14 +457,17 @@ def normalise_operation(operation: str, payload: dict[str, Any]) -> tuple[dict[s
|
||||
"files": files, "checks": checks, "compose_file": compose_file,
|
||||
"services": [str(x) for x in services], "build": bool(payload.get("build", True)),
|
||||
"openwebui_sync": bool(payload.get("openwebui_sync", True)),
|
||||
"hermes_sync": bool(payload.get("hermes_sync", False)),
|
||||
"message": message, "paths": selected,
|
||||
"create_recovery": bool(payload.get("create_recovery", True)), "recovery_label": label,
|
||||
}
|
||||
preview = (
|
||||
f"ONE MCP RELEASE\nServices: {', '.join(normal['services'])}\n"
|
||||
f"Checks: {', '.join(checks)}\nOpenWebUI sync: {normal['openwebui_sync']}\n"
|
||||
f"Hermes sync: {normal['hermes_sync']}\n"
|
||||
f"Selective commit: {message}\nPaths: {', '.join(selected)}\n"
|
||||
f"Recovery: {normal['create_recovery']} ({label})\n\n{patch_preview}"
|
||||
f"Recovery: {normal['create_recovery']} ({label})\n\n"
|
||||
+ "\n".join(part for part in (import_preview, patch_preview) if part)
|
||||
)
|
||||
return normal, preview
|
||||
if operation == "run_checks":
|
||||
@@ -579,6 +653,7 @@ def execute_operation(ticket: str, operation: str, payload: dict[str, Any]) -> d
|
||||
published = False
|
||||
deployed = False
|
||||
synced = False
|
||||
hermes_synced = False
|
||||
try:
|
||||
changed = apply_files(payload["files"], backup)
|
||||
checks = []
|
||||
@@ -595,10 +670,14 @@ def execute_operation(ticket: str, operation: str, payload: dict[str, Any]) -> d
|
||||
if payload["openwebui_sync"]:
|
||||
sync = run(["bash", str(STACK / "platform/openwebui/install-filters.sh")], cwd=STACK, timeout=1800, check=True)
|
||||
synced = True
|
||||
hermes_sync = None
|
||||
if payload["hermes_sync"]:
|
||||
hermes_sync = run(["bash", str(STACK / "platform/hermes/install-hermes.sh")], cwd=STACK, timeout=1800, check=True)
|
||||
hermes_synced = True
|
||||
publication = publish_paths(payload["message"], payload["paths"])
|
||||
published = True
|
||||
recovery = perform_recovery(payload["recovery_label"]) if payload["create_recovery"] else None
|
||||
return {"changed": changed, "checks": checks, "deploy": deploy, "openwebui_sync": sync, "publication": publication, "recovery": recovery, "containers": run(["docker", "ps", "--format", "{{.Names}}\t{{.Status}}"])}
|
||||
return {"changed": changed, "checks": checks, "deploy": deploy, "openwebui_sync": sync, "hermes_sync": hermes_sync, "publication": publication, "recovery": recovery, "containers": run(["docker", "ps", "--format", "{{.Names}}\t{{.Status}}"])}
|
||||
except Exception:
|
||||
if not published:
|
||||
restore_files(payload["files"], backup)
|
||||
@@ -608,6 +687,8 @@ def execute_operation(ticket: str, operation: str, payload: dict[str, Any]) -> d
|
||||
run(rollback_argv, cwd=STACK, timeout=3600)
|
||||
if synced:
|
||||
run(["bash", str(STACK / "platform/openwebui/install-filters.sh")], cwd=STACK, timeout=1800)
|
||||
if hermes_synced:
|
||||
run(["bash", str(STACK / "platform/hermes/install-hermes.sh")], cwd=STACK, timeout=1800)
|
||||
raise
|
||||
return start_job(ticket, operation, release)
|
||||
if operation == "run_checks":
|
||||
|
||||
Reference in New Issue
Block a user