Add compact Athena MCP release workflow
This commit is contained in:
@@ -28,7 +28,7 @@ from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
|
||||
VERSION = "2.0.0"
|
||||
VERSION = "2.2.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()
|
||||
@@ -48,9 +48,11 @@ PROTECTED_CONTAINERS = {
|
||||
"mike-ai-wireguard-gateway",
|
||||
}
|
||||
ALLOWED_OPERATIONS = {
|
||||
"file_update", "run_checks", "compose_deploy", "container_action",
|
||||
"file_update", "patch_update", "mcp_release", "run_checks", "compose_deploy", "container_action",
|
||||
"openwebui_sync", "git_publish", "model_download", "benchmark", "recovery",
|
||||
}
|
||||
|
||||
HUNK_HEADER = re.compile(r"^@@ -(\d+)(?:,(\d+))? \+(\d+)(?:,(\d+))? @@")
|
||||
ALLOWED_CHECKS = {
|
||||
"operator-tests": ["python3", "dev/test_athena_operator.py"],
|
||||
"openwebui-filter-tests": ["python3", "dev/test_openwebui_filters.py"],
|
||||
@@ -134,6 +136,85 @@ def sha(data: bytes) -> str:
|
||||
return hashlib.sha256(data).hexdigest()
|
||||
|
||||
|
||||
def apply_unified_patch(before: str, patch: str) -> str:
|
||||
"""Apply one ordinary unified diff without invoking a shell command."""
|
||||
if not patch or len(patch.encode()) > MAX_FILE_BYTES:
|
||||
raise ValueError("patch is empty or exceeds limit")
|
||||
source = before.splitlines(keepends=True)
|
||||
lines = patch.splitlines(keepends=True)
|
||||
position = 0
|
||||
output: list[str] = []
|
||||
index = 0
|
||||
if index < len(lines) and lines[index].startswith("--- "):
|
||||
index += 1
|
||||
if index < len(lines) and lines[index].startswith("+++ "):
|
||||
index += 1
|
||||
hunks = 0
|
||||
while index < len(lines):
|
||||
header = lines[index].rstrip("\r\n")
|
||||
match = HUNK_HEADER.match(header)
|
||||
if not match:
|
||||
raise ValueError("patch must contain only standard unified-diff hunks")
|
||||
old_start = int(match.group(1))
|
||||
old_count = int(match.group(2) or "1")
|
||||
new_count = int(match.group(4) or "1")
|
||||
target = max(0, old_start - 1)
|
||||
if target < position or target > len(source):
|
||||
raise RuntimeError("patch hunk is outside the source file")
|
||||
output.extend(source[position:target])
|
||||
position = target
|
||||
consumed = produced = 0
|
||||
index += 1
|
||||
while index < len(lines) and not lines[index].startswith("@@ "):
|
||||
line = lines[index]
|
||||
if line.startswith("\\ No newline at end of file"):
|
||||
index += 1
|
||||
continue
|
||||
if not line or line[0] not in " +-":
|
||||
raise ValueError("invalid unified-diff line")
|
||||
marker, value = line[0], line[1:]
|
||||
if marker in " -":
|
||||
if position >= len(source) or source[position] != value:
|
||||
raise RuntimeError("patch context does not match source")
|
||||
if marker == " ":
|
||||
output.append(source[position]); produced += 1
|
||||
position += 1; consumed += 1
|
||||
else:
|
||||
output.append(value); produced += 1
|
||||
index += 1
|
||||
if consumed != old_count or produced != new_count:
|
||||
raise RuntimeError("patch hunk line counts do not match its header")
|
||||
hunks += 1
|
||||
if not hunks:
|
||||
raise ValueError("patch contains no hunks")
|
||||
output.extend(source[position:])
|
||||
return "".join(output)
|
||||
|
||||
|
||||
def normalise_patches(files: Any) -> tuple[list[dict[str, Any]], str]:
|
||||
if not isinstance(files, list) or not 1 <= len(files) <= MAX_FILES:
|
||||
raise ValueError("files must contain 1..24 patch entries")
|
||||
normal: list[dict[str, Any]] = []
|
||||
previews: list[str] = []
|
||||
for item in files:
|
||||
if not isinstance(item, dict):
|
||||
raise ValueError("each patch entry must be an object")
|
||||
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 = str(item.get("expected_sha256", ""))
|
||||
if expected and expected != before_sha:
|
||||
raise RuntimeError(f"source drift for {relative}")
|
||||
patch = str(item.get("patch", ""))
|
||||
after = apply_unified_patch(before, patch)
|
||||
if len(after.encode()) > MAX_FILE_BYTES:
|
||||
raise ValueError("patched file exceeds limit")
|
||||
normal.append({"path": str(relative), "content": after, "before_sha256": before_sha, "after_sha256": sha(after.encode())})
|
||||
previews.append(f"### {relative}\n{compact(patch)}")
|
||||
return normal, "\n\n".join(previews)
|
||||
|
||||
|
||||
def json_write(path: Path, value: Any) -> None:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
temporary = path.with_suffix(path.suffix + ".tmp")
|
||||
@@ -277,6 +358,44 @@ def normalise_operation(operation: str, payload: dict[str, Any]) -> tuple[dict[s
|
||||
normal.append({"path": str(relative), "content": content, "before_sha256": before_sha, "after_sha256": sha(raw)})
|
||||
previews.append(compact(diff))
|
||||
return {"files": normal}, "\n".join(previews)
|
||||
if operation == "patch_update":
|
||||
files, preview = normalise_patches(payload.get("files"))
|
||||
return {"files": files}, preview
|
||||
if operation == "mcp_release":
|
||||
files, patch_preview = normalise_patches(payload.get("files"))
|
||||
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")
|
||||
compose_file = str(payload.get("compose_file", "platform/mcp/compose.yaml"))
|
||||
if compose_file != "platform/mcp/compose.yaml":
|
||||
raise ValueError("MCP releases must use platform/mcp/compose.yaml")
|
||||
services = payload.get("services") or []
|
||||
if not isinstance(services, list) or not 1 <= len(services) <= 12 or any(not SAFE_NAME.fullmatch(str(x)) for x in services):
|
||||
raise ValueError("invalid MCP service list")
|
||||
message = str(payload.get("message", ""))
|
||||
if not SAFE_COMMIT.fullmatch(message):
|
||||
raise ValueError("invalid commit message")
|
||||
paths = payload.get("paths") or [item["path"] for item in files]
|
||||
selected = [str(safe_relative(str(path))) for path in paths]
|
||||
if len(set(selected)) != len(selected) or not set(item["path"] for item in files).issubset(set(selected)):
|
||||
raise ValueError("release paths must be unique and include every patched file")
|
||||
label = str(payload.get("recovery_label", time.strftime("%Y%m%d-%H%M")))
|
||||
if not SAFE_NAME.fullmatch(label):
|
||||
raise ValueError("invalid recovery label")
|
||||
normal = {
|
||||
"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)),
|
||||
"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"Selective commit: {message}\nPaths: {', '.join(selected)}\n"
|
||||
f"Recovery: {normal['create_recovery']} ({label})\n\n{patch_preview}"
|
||||
)
|
||||
return normal, preview
|
||||
if operation == "run_checks":
|
||||
checks = payload.get("checks") or []
|
||||
if not isinstance(checks, list) or not checks or any(name not in ALLOWED_CHECKS for name in checks):
|
||||
@@ -390,18 +509,107 @@ def start_job(ticket: str, operation: str, worker) -> dict[str, Any]:
|
||||
return {"job_id": job_id, "status": "running", "instruction": "Poll athena_operator_job until completed or failed."}
|
||||
|
||||
|
||||
def apply_files(files: list[dict[str, Any]], backup: Path) -> list[str]:
|
||||
for item in files:
|
||||
relative = safe_relative(item["path"])
|
||||
current = source_file(REPOSITORY, relative)
|
||||
current_sha = sha(current.read_bytes()) if current.is_file() else sha(b"")
|
||||
if current_sha != item["before_sha256"]:
|
||||
raise RuntimeError(f"source drift after preview: {relative}")
|
||||
for item in files:
|
||||
sync_file(safe_relative(item["path"]), item["content"], backup)
|
||||
return [item["path"] for item in files]
|
||||
|
||||
|
||||
def restore_files(files: list[dict[str, Any]], backup: Path) -> None:
|
||||
for item in files:
|
||||
relative = safe_relative(item["path"])
|
||||
for root in (REPOSITORY, STACK):
|
||||
target = source_file(root, relative)
|
||||
saved = backup / root.name / relative
|
||||
if saved.is_file():
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
shutil.copy2(saved, target)
|
||||
elif target.exists():
|
||||
target.unlink()
|
||||
|
||||
|
||||
def publish_paths(message: str, selected: list[str]) -> dict[str, Any]:
|
||||
run(["git", "add", "--", *selected], cwd=REPOSITORY, check=True)
|
||||
run(["git", "diff", "--cached", "--check", "--", *selected], cwd=REPOSITORY, check=True)
|
||||
commit = run(["git", "commit", "-m", message, "--", *selected], cwd=REPOSITORY, check=True)
|
||||
pushed = run(["git", "push", "origin", "HEAD:main"], cwd=REPOSITORY, timeout=300, check=True)
|
||||
head = run(["git", "rev-parse", "HEAD"], cwd=REPOSITORY, check=True)["output"].strip()
|
||||
(STACK / ".mike-ai-source-commit").write_text(head + "\n")
|
||||
return {"commit": head, "commit_output": commit, "push_output": pushed}
|
||||
|
||||
|
||||
def perform_recovery(label: str) -> dict[str, Any]:
|
||||
dirty = run(["git", "status", "--porcelain"], cwd=REPOSITORY, check=True)["output"].strip()
|
||||
if dirty:
|
||||
raise RuntimeError("publish repository changes before creating a recovery kit")
|
||||
head = run(["git", "rev-parse", "HEAD"], cwd=REPOSITORY, check=True)["output"].strip()
|
||||
marker = (STACK / ".mike-ai-source-commit").read_text().strip()
|
||||
if marker != head:
|
||||
raise RuntimeError("deployed source marker and operator repository HEAD differ")
|
||||
encrypted = Path(f"/data/athena-recovery-{label}.tar.age")
|
||||
source_bundle = Path(f"/data/athena-source-{label}.git.bundle")
|
||||
release = Path(f"/data/mike-ai-recovery-kit-{label}")
|
||||
for target in (encrypted, source_bundle, release):
|
||||
if target.exists():
|
||||
raise FileExistsError(target)
|
||||
git_bundle = run(["git", "bundle", "create", str(source_bundle), "--all"], cwd=REPOSITORY, timeout=1800, check=True)
|
||||
encrypted_result = run([str(STACK / "platform/recovery/create-recovery-bundle.sh"), str(encrypted)], timeout=7200, check=True)
|
||||
identity = Path("/data/mike-ai-recovery-kit/recovery.agekey")
|
||||
if not identity.is_file():
|
||||
raise RuntimeError("existing recovery identity is unavailable")
|
||||
kit_result = run([str(STACK / "platform/recovery/create-self-contained-data-kit.sh"), str(encrypted), str(identity), str(source_bundle), str(release)], timeout=7200, check=True)
|
||||
verify = run(["sha256sum", "-c", "SHA256SUMS"], cwd=release, timeout=1800, check=True)
|
||||
return {"commit": head, "recovery_bundle": str(encrypted), "source_bundle": str(source_bundle), "self_contained_kit": str(release), "git_bundle": git_bundle, "encrypted_bundle": encrypted_result, "kit": kit_result, "verification": verify}
|
||||
|
||||
|
||||
def execute_operation(ticket: str, operation: str, payload: dict[str, Any]) -> dict[str, Any]:
|
||||
if operation == "file_update":
|
||||
if operation in {"file_update", "patch_update"}:
|
||||
backup = STATE / "backups" / f"{now()}-{ticket}"
|
||||
for item in payload["files"]:
|
||||
relative = safe_relative(item["path"])
|
||||
current = source_file(REPOSITORY, relative)
|
||||
current_sha = sha(current.read_bytes()) if current.is_file() else sha(b"")
|
||||
if current_sha != item["before_sha256"]:
|
||||
raise RuntimeError(f"source drift after preview: {relative}")
|
||||
for item in payload["files"]:
|
||||
sync_file(safe_relative(item["path"]), item["content"], backup)
|
||||
return {"changed": [item["path"] for item in payload["files"]], "backup": str(backup), "git_diff": run(["git", "diff", "--stat"], cwd=REPOSITORY)}
|
||||
changed = apply_files(payload["files"], backup)
|
||||
return {"changed": changed, "backup": str(backup), "git_diff": run(["git", "diff", "--stat"], cwd=REPOSITORY)}
|
||||
if operation == "mcp_release":
|
||||
def release():
|
||||
backup = STATE / "backups" / f"{now()}-{ticket}"
|
||||
published = False
|
||||
deployed = False
|
||||
synced = False
|
||||
try:
|
||||
changed = apply_files(payload["files"], backup)
|
||||
checks = []
|
||||
for name in payload["checks"]:
|
||||
result = run(ALLOWED_CHECKS[name], cwd=REPOSITORY, timeout=1200, check=True)
|
||||
checks.append({"name": name, **result})
|
||||
run(["docker", "compose", "-f", payload["compose_file"], "config", "-q"], cwd=STACK, check=True)
|
||||
argv = ["docker", "compose", "-f", payload["compose_file"], "up", "-d"]
|
||||
if payload["build"]: argv.append("--build")
|
||||
argv.extend(payload["services"])
|
||||
deploy = run(argv, cwd=STACK, timeout=3600, check=True)
|
||||
deployed = True
|
||||
sync = None
|
||||
if payload["openwebui_sync"]:
|
||||
sync = run(["bash", str(STACK / "platform/openwebui/install-filters.sh")], cwd=STACK, timeout=1800, check=True)
|
||||
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}}"])}
|
||||
except Exception:
|
||||
if not published:
|
||||
restore_files(payload["files"], backup)
|
||||
run(["git", "reset", "--", *payload["paths"]], cwd=REPOSITORY)
|
||||
if deployed:
|
||||
rollback_argv = ["docker", "compose", "-f", payload["compose_file"], "up", "-d", "--build", *payload["services"]]
|
||||
run(rollback_argv, cwd=STACK, timeout=3600)
|
||||
if synced:
|
||||
run(["bash", str(STACK / "platform/openwebui/install-filters.sh")], cwd=STACK, timeout=1800)
|
||||
raise
|
||||
return start_job(ticket, operation, release)
|
||||
if operation == "run_checks":
|
||||
return {"checks": [{"name": name, **run(ALLOWED_CHECKS[name], cwd=REPOSITORY, timeout=1200)} for name in payload["checks"]]}
|
||||
if operation == "compose_deploy":
|
||||
@@ -424,13 +632,7 @@ def execute_operation(ticket: str, operation: str, payload: dict[str, Any]) -> d
|
||||
current_status = run(["git", "status", "--short", "--", *selected], cwd=REPOSITORY, check=True)["output"].strip()
|
||||
if current_status != payload["reviewed_status"]:
|
||||
raise RuntimeError("selected repository paths changed after the Git publish preview")
|
||||
run(["git", "add", "--", *selected], cwd=REPOSITORY, check=True)
|
||||
run(["git", "diff", "--cached", "--check", "--", *selected], cwd=REPOSITORY, check=True)
|
||||
commit = run(["git", "commit", "-m", payload["message"], "--", *selected], cwd=REPOSITORY, check=True)
|
||||
pushed = run(["git", "push", "origin", "HEAD:main"], cwd=REPOSITORY, timeout=300, check=True)
|
||||
head = run(["git", "rev-parse", "HEAD"], cwd=REPOSITORY, check=True)["output"].strip()
|
||||
(STACK / ".mike-ai-source-commit").write_text(head + "\n")
|
||||
return {"commit": head, "commit_output": commit, "push_output": pushed}
|
||||
return publish_paths(payload["message"], selected)
|
||||
if operation == "model_download":
|
||||
def download():
|
||||
destination = source_file(MODELS, Path(payload["destination"]))
|
||||
@@ -456,40 +658,7 @@ def execute_operation(ticket: str, operation: str, payload: dict[str, Any]) -> d
|
||||
return start_job(ticket, operation, benchmark)
|
||||
if operation == "recovery":
|
||||
def recovery():
|
||||
label = payload["label"]
|
||||
dirty = run(["git", "status", "--porcelain"], cwd=REPOSITORY, check=True)["output"].strip()
|
||||
if dirty:
|
||||
raise RuntimeError("publish repository changes before creating a recovery kit")
|
||||
head = run(["git", "rev-parse", "HEAD"], cwd=REPOSITORY, check=True)["output"].strip()
|
||||
marker = (STACK / ".mike-ai-source-commit").read_text().strip()
|
||||
if marker != head:
|
||||
raise RuntimeError("deployed source marker and operator repository HEAD differ")
|
||||
encrypted = Path(f"/data/athena-recovery-{label}.tar.age")
|
||||
source_bundle = Path(f"/data/athena-source-{label}.git.bundle")
|
||||
release = Path(f"/data/mike-ai-recovery-kit-{label}")
|
||||
for target in (encrypted, source_bundle, release):
|
||||
if target.exists():
|
||||
raise FileExistsError(target)
|
||||
git_bundle = run(["git", "bundle", "create", str(source_bundle), "--all"], cwd=REPOSITORY, timeout=1800, check=True)
|
||||
encrypted_result = run([str(STACK / "platform/recovery/create-recovery-bundle.sh"), str(encrypted)], timeout=7200, check=True)
|
||||
identity = Path("/data/mike-ai-recovery-kit/recovery.agekey")
|
||||
if not identity.is_file():
|
||||
raise RuntimeError("existing recovery identity is unavailable")
|
||||
kit_result = run([
|
||||
str(STACK / "platform/recovery/create-self-contained-data-kit.sh"),
|
||||
str(encrypted), str(identity), str(source_bundle), str(release),
|
||||
], timeout=7200, check=True)
|
||||
verify = run(["sha256sum", "-c", "SHA256SUMS"], cwd=release, timeout=1800, check=True)
|
||||
return {
|
||||
"commit": head,
|
||||
"recovery_bundle": str(encrypted),
|
||||
"source_bundle": str(source_bundle),
|
||||
"self_contained_kit": str(release),
|
||||
"git_bundle": git_bundle,
|
||||
"encrypted_bundle": encrypted_result,
|
||||
"kit": kit_result,
|
||||
"verification": verify,
|
||||
}
|
||||
return perform_recovery(payload["label"])
|
||||
return start_job(ticket, operation, recovery)
|
||||
raise AssertionError(operation)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user