feat(poc): rewind MCP, skill, and recover-ladder docs

Initial PoC: find_protection, tagged checkpoints, FLR, mutating catalog.
Lab 10.9 status enums (0=Initializing, 1=MeetingSLA). Credentials stay in gitignored config.json.
This commit is contained in:
2026-09-21 12:09:11 -04:00
commit 38ba9c1b50
24 changed files with 1811 additions and 0 deletions
+3
View File
@@ -0,0 +1,3 @@
"""Zerto AI Rewind PoC MCP."""
__version__ = "0.1.0"
+4
View File
@@ -0,0 +1,4 @@
from zerto_rewind_mcp.server import main
if __name__ == "__main__":
main()
+73
View File
@@ -0,0 +1,73 @@
"""Mutating catalog: opt-in list of tools that must be guarded."""
from __future__ import annotations
import json
from dataclasses import asdict, dataclass
from pathlib import Path
from typing import Any
@dataclass
class CatalogEntry:
server: str
tool: str
vm_arg: str
notes: str = ""
def key(self) -> tuple[str, str]:
return (self.server.lower(), self.tool.lower())
def as_dict(self) -> dict[str, str]:
return asdict(self)
def entry_from_dict(raw: dict[str, Any]) -> CatalogEntry:
server = str(raw.get("server") or "").strip()
tool = str(raw.get("tool") or "").strip()
vm_arg = str(raw.get("vm_arg") or "").strip()
if not server or not tool or not vm_arg:
raise ValueError("catalog entry needs server, tool, and vm_arg")
return CatalogEntry(
server=server,
tool=tool,
vm_arg=vm_arg,
notes=str(raw.get("notes") or ""),
)
class MutatingCatalog:
def __init__(self, entries: list[CatalogEntry] | None = None, path: Path | None = None):
self._entries: dict[tuple[str, str], CatalogEntry] = {}
self.path = path
for entry in entries or []:
self._entries[entry.key()] = entry
def list(self) -> list[CatalogEntry]:
return sorted(self._entries.values(), key=lambda e: e.key())
def get(self, server: str, tool: str) -> CatalogEntry | None:
return self._entries.get((server.lower(), tool.lower()))
def add(self, entry: CatalogEntry) -> CatalogEntry:
self._entries[entry.key()] = entry
self.save()
return entry
def save(self) -> None:
if self.path is None:
return
payload = {"mutating_tools": [e.as_dict() for e in self.list()]}
self.path.parent.mkdir(parents=True, exist_ok=True)
if self.path.suffix == ".json" and self.path.exists():
existing = json.loads(self.path.read_text(encoding="utf-8"))
if isinstance(existing, dict):
existing["mutating_tools"] = payload["mutating_tools"]
payload = existing
self.path.write_text(json.dumps(payload, indent=2) + "\n", encoding="utf-8")
@classmethod
def from_config(cls, data: dict[str, Any], path: Path | None = None) -> MutatingCatalog:
raw = data.get("mutating_tools") or []
entries = [entry_from_dict(item) for item in raw]
return cls(entries, path=path)
+105
View File
@@ -0,0 +1,105 @@
"""Tagged checkpoint insert and wait-until-listed."""
from __future__ import annotations
import asyncio
import re
from datetime import UTC, datetime
from typing import Any
from zerto_rewind_mcp.client import ZertoClient, ZertoError
from zerto_rewind_mcp.protection import FindResult
from zerto_rewind_mcp.util import pick
_SAFE = re.compile(r"[^A-Za-z0-9._:-]+")
def make_tag(agent: str, change_id: str, when: datetime | None = None) -> str:
stamp = (when or datetime.now(UTC)).strftime("%Y%m%dT%H%M%SZ")
agent_s = _SAFE.sub("-", agent.strip())[:40] or "agent"
change_s = _SAFE.sub("-", change_id.strip())[:80] or "change"
return f"ai:{agent_s}:{change_s}:{stamp}"
def checkpoint_tag(row: dict[str, Any]) -> str:
value = pick(row, "Tag", "tag", "CheckpointName", "checkpointName") or ""
return str(value)
def checkpoint_id(row: dict[str, Any]) -> str:
value = pick(
row, "CheckpointId", "checkpointId", "checkpointIdentifier", "CheckpointIdentifier"
)
return str(value) if value is not None else ""
async def wait_for_tag(
client: ZertoClient,
vpg_identifier: str,
tag: str,
*,
timeout_s: float = 45.0,
interval_s: float = 1.5,
) -> dict[str, Any]:
deadline = asyncio.get_event_loop().time() + timeout_s
last: list[dict[str, Any]] = []
while asyncio.get_event_loop().time() < deadline:
last = await client.list_checkpoints(vpg_identifier)
for row in last:
if checkpoint_tag(row) == tag:
return row
await asyncio.sleep(interval_s)
raise ZertoError(
f"Tagged checkpoint {tag!r} did not appear on VPG {vpg_identifier} "
f"within {timeout_s:.0f}s. Do not mutate. "
"If the protected site is Azure or AWS, tagged checkpoints are not supported."
)
async def tag_vpgs(
client: ZertoClient,
result: FindResult,
tag: str,
) -> dict[str, Any]:
if result.outcome != "ok" or result.vm is None:
return {
"ok": False,
"tag": tag,
"message": result.message,
"find": result.as_dict(),
}
tagged: list[dict[str, Any]] = []
skipped = [v.as_dict() for v in result.vm.vpgs if not v.can_tag]
errors: list[str] = []
for vpg in result.taggable_vpgs:
try:
insert = await client.insert_checkpoint(vpg.vpg_identifier, tag)
row = await wait_for_tag(client, vpg.vpg_identifier, tag)
tagged.append(
{
"vpg_identifier": vpg.vpg_identifier,
"vpg_name": vpg.vpg_name,
"checkpoint_id": checkpoint_id(row),
"tag": checkpoint_tag(row) or tag,
"insert_result": insert if not isinstance(insert, dict) else "ok",
}
)
except ZertoError as exc:
errors.append(f"{vpg.vpg_name} ({vpg.vpg_identifier}): {exc}")
ok = bool(tagged) and not errors
message = (
f"Tagged {len(tagged)} VPG(s) with {tag!r}."
if ok
else f"Checkpoint failed on {len(errors)} VPG(s). Refuse the change. " + "; ".join(errors)
)
if skipped and ok:
message += " Some VPGs were skipped (syncing or not Protecting)."
return {
"ok": ok,
"tag": tag,
"vm": result.vm.as_dict(),
"tagged": tagged,
"skipped": skipped,
"errors": errors,
"message": message,
}
+276
View File
@@ -0,0 +1,276 @@
"""ZVM / ZCA REST client. Same paths on both (zerto_api_lessons)."""
from __future__ import annotations
import time
from typing import Any
from urllib.parse import urljoin
import httpx
from zerto_rewind_mcp.util import pick
class ZertoError(Exception):
def __init__(self, message: str, status_code: int | None = None, body: str | None = None):
super().__init__(message)
self.status_code = status_code
self.body = body
class ZertoClient:
def __init__(
self,
base_url: str,
username: str,
password: str,
client_id: str = "zerto-client",
verify_tls: bool = False,
timeout: float = 60.0,
):
self.base_url = base_url.rstrip("/")
self.username = username
self.password = password
self.client_id = client_id
self._token: str | None = None
self._token_exp = 0.0
self._http = httpx.AsyncClient(verify=verify_tls, timeout=timeout)
async def aclose(self) -> None:
await self._http.aclose()
def _url(self, path: str) -> str:
if path.startswith("http"):
return path
return urljoin(self.base_url + "/", path.lstrip("/"))
async def ensure_token(self) -> str:
if self._token and time.time() < self._token_exp - 30:
return self._token
url = self._url("/auth/realms/zerto/protocol/openid-connect/token")
response = await self._http.post(
url,
data={
"grant_type": "password",
"username": self.username,
"password": self.password,
"client_id": self.client_id,
"scope": "openid",
},
headers={"Content-Type": "application/x-www-form-urlencoded"},
)
if response.status_code >= 400:
raise ZertoError(
f"Keycloak token failed HTTP {response.status_code}. "
"Check username/password/client_id "
"(10.x uses zerto-client; 9.x may use zerto-api).",
status_code=response.status_code,
body=response.text[:500],
)
payload = response.json()
self._token = payload["access_token"]
self._token_exp = time.time() + float(payload.get("expires_in") or 60)
return self._token
async def request(
self,
method: str,
path: str,
*,
params: dict[str, Any] | None = None,
json_body: Any = None,
) -> httpx.Response:
token = await self.ensure_token()
headers = {"Authorization": f"Bearer {token}"}
if json_body is not None:
headers["Content-Type"] = "application/json"
response = await self._http.request(
method,
self._url(path),
params=params,
json=json_body,
headers=headers,
)
if response.status_code == 401:
self._token = None
token = await self.ensure_token()
headers["Authorization"] = f"Bearer {token}"
response = await self._http.request(
method,
self._url(path),
params=params,
json=json_body,
headers=headers,
)
return response
async def json(
self,
method: str,
path: str,
*,
params: dict[str, Any] | None = None,
json_body: Any = None,
) -> Any:
response = await self.request(method, path, params=params, json_body=json_body)
if response.status_code >= 400:
raise ZertoError(
f"{method} {path} HTTP {response.status_code}: {response.text[:400]}",
status_code=response.status_code,
body=response.text[:1000],
)
if response.status_code == 204 or not response.content:
return None
try:
return response.json()
except ValueError:
return response.text
async def get_vms(
self,
*,
vm_name: str | None = None,
vm_identifier: str | None = None,
) -> list[dict[str, Any]]:
params: dict[str, Any] = {}
if vm_identifier:
params["vmIdentifier"] = vm_identifier
if vm_name:
params["vmName"] = vm_name
data = await self.json("GET", "/v1/vms", params=params or None)
if data is None:
return []
if isinstance(data, list):
return data
if isinstance(data, dict):
inner = pick(data, "value", "items", "vms")
if isinstance(inner, list):
return inner
raise ZertoError(f"GET /v1/vms returned unexpected shape: {type(data).__name__}")
async def get_vpg(self, vpg_identifier: str) -> dict[str, Any]:
data = await self.json("GET", f"/v1/vpgs/{vpg_identifier}")
if not isinstance(data, dict):
raise ZertoError("GET /v1/vpgs/{id} did not return an object")
return data
async def list_checkpoints(self, vpg_identifier: str) -> list[dict[str, Any]]:
data = await self.json("GET", f"/v1/vpgs/{vpg_identifier}/checkpoints")
if data is None:
return []
if isinstance(data, list):
return data
raise ZertoError("GET checkpoints returned unexpected shape")
async def insert_checkpoint(self, vpg_identifier: str, tag: str) -> Any:
"""POST tagged checkpoint. Body key is CheckpointName on documented 9.x API."""
path = f"/v1/vpgs/{vpg_identifier}/checkpoints"
try:
return await self.json("POST", path, json_body={"CheckpointName": tag})
except ZertoError as exc:
if exc.status_code != 400:
raise
return await self.json("POST", path, json_body={"checkpointName": tag})
async def get_task(self, task_id: str) -> dict[str, Any]:
data = await self.json("GET", f"/v1/tasks/{task_id}")
if not isinstance(data, dict):
raise ZertoError("GET /v1/tasks/{id} did not return an object")
return data
async def start_flr(
self,
vpg_identifier: str,
vm_identifier: str,
checkpoint_identifier: str,
initial_download_path: str,
) -> Any:
return await self.json(
"POST",
"/v1/flrs",
json_body={
"jflr": {
"vpgIdentifier": vpg_identifier,
"vmIdentifier": vm_identifier,
"CheckpointIdentifier": checkpoint_identifier,
"initialDownloadPath": initial_download_path,
}
},
)
async def get_flr(self, session_id: str) -> Any:
return await self.json("GET", f"/v1/flrs/{session_id}")
async def browse_flr(self, session_id: str, path: str = "", recursive: bool = False) -> Any:
return await self.json(
"POST",
f"/v1/flrs/{session_id}/browse",
json_body={"path": path, "recursive": recursive},
)
async def download_flr(self, session_id: str, path_list: list[str]) -> Any:
return await self.json(
"POST",
f"/v1/flrs/{session_id}/download",
json_body={"pathList": path_list},
)
async def fetch_download(self, token: str) -> bytes:
response = await self.request("GET", f"/v1/downloads/{token}")
if response.status_code >= 400:
# workflow also shows GET /v1/flrs/{token}
response = await self.request("GET", f"/v1/flrs/{token}")
if response.status_code >= 400:
raise ZertoError(
f"FLR download HTTP {response.status_code}: {response.text[:300]}",
status_code=response.status_code,
)
return response.content
async def end_flr(self, session_id: str) -> None:
try:
await self.json("DELETE", f"/v1/flrs/{session_id}")
except ZertoError:
await self.json("POST", f"/v1/flrs/{session_id}")
async def start_failover_test(
self,
vpg_identifier: str,
checkpoint_identifier: str | None = None,
vm_identifiers: list[str] | None = None,
) -> Any:
body: dict[str, Any] = {}
if checkpoint_identifier:
body["checkpointIdentifier"] = checkpoint_identifier
if vm_identifiers:
body["vmIdentifiers"] = vm_identifiers
return await self.json(
"POST",
f"/v1/vpgs/{vpg_identifier}/FailoverTest",
json_body=body or None,
)
async def stop_failover_test(self, vpg_identifier: str, success: bool, summary: str) -> Any:
return await self.json(
"POST",
f"/v1/vpgs/{vpg_identifier}/FailoverTestStop",
json_body={"failoverTestSuccess": success, "failoverTestSummary": summary},
)
async def start_clone(
self,
vpg_identifier: str,
checkpoint_identifier: str,
datastore_identifier: str | None = None,
vm_identifiers: list[str] | None = None,
) -> Any:
body: dict[str, Any] = {"checkpointIdentifier": checkpoint_identifier}
if datastore_identifier:
body["datastoreIdentifier"] = datastore_identifier
if vm_identifiers:
body["vmIdentifiers"] = vm_identifiers
return await self.json(
"POST",
f"/v1/vpgs/{vpg_identifier}/CloneStart",
json_body=body,
)
+42
View File
@@ -0,0 +1,42 @@
"""Load config from JSON and/or environment."""
from __future__ import annotations
import json
import os
from pathlib import Path
from typing import Any
from zerto_rewind_mcp.catalog import MutatingCatalog
def _env(name: str, default: str | None = None) -> str | None:
value = os.environ.get(name)
if value is None or value == "":
return default
return value
def load_config(path: str | Path | None = None) -> dict[str, Any]:
config_path = Path(path or _env("ZERTO_REWIND_CONFIG") or "config.json")
data: dict[str, Any] = {}
if config_path.exists():
data = json.loads(config_path.read_text(encoding="utf-8"))
data["_config_path"] = str(config_path.resolve())
data.setdefault("zerto_url", _env("ZERTO_URL", ""))
data.setdefault("username", _env("ZERTO_USERNAME", ""))
data.setdefault("password", _env("ZERTO_PASSWORD", ""))
data.setdefault("client_id", _env("ZERTO_CLIENT_ID", "zerto-client"))
verify = _env("ZERTO_VERIFY_TLS")
if verify is not None:
data["verify_tls"] = verify.lower() not in {"0", "false", "no"}
else:
data.setdefault("verify_tls", False)
data.setdefault("recovery_dir", _env("ZERTO_RECOVERY_DIR", "./recovered"))
data.setdefault("mutating_tools", data.get("mutating_tools") or [])
return data
def load_catalog(data: dict[str, Any]) -> MutatingCatalog:
path = Path(data["_config_path"]) if data.get("_config_path") else None
return MutatingCatalog.from_config(data, path=path)
+170
View File
@@ -0,0 +1,170 @@
"""find_protection: unique VM + every VPG, or stop."""
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any, Literal
from zerto_rewind_mcp.status import can_tag, status_name, substatus_name
from zerto_rewind_mcp.util import pick
Outcome = Literal["none", "ambiguous", "ok"]
@dataclass
class VpgMembership:
vpg_identifier: str
vpg_name: str
status: int | None
sub_status: int | None
can_tag: bool
skip_reason: str | None
def as_dict(self) -> dict[str, Any]:
return {
"vpg_identifier": self.vpg_identifier,
"vpg_name": self.vpg_name,
"status": status_name(self.status),
"sub_status": substatus_name(self.sub_status),
"can_tag": self.can_tag,
"skip_reason": self.skip_reason,
}
@dataclass
class VmMatch:
vm_identifier: str
vm_name: str
vpgs: list[VpgMembership] = field(default_factory=list)
def as_dict(self) -> dict[str, Any]:
return {
"vm_identifier": self.vm_identifier,
"vm_name": self.vm_name,
"vpgs": [v.as_dict() for v in self.vpgs],
}
@dataclass
class FindResult:
outcome: Outcome
query: str
message: str
vm: VmMatch | None = None
matches: list[VmMatch] = field(default_factory=list)
def as_dict(self) -> dict[str, Any]:
payload: dict[str, Any] = {
"outcome": self.outcome,
"query": self.query,
"message": self.message,
}
if self.vm is not None:
payload["vm"] = self.vm.as_dict()
if self.matches:
payload["matches"] = [m.as_dict() for m in self.matches]
return payload
@property
def taggable_vpgs(self) -> list[VpgMembership]:
if self.vm is None:
return []
return [v for v in self.vm.vpgs if v.can_tag]
def membership_from_row(row: dict[str, Any]) -> VpgMembership:
status = pick(row, "Status", "status")
sub = pick(row, "SubStatus", "subStatus", "sub_status")
ok, reason = can_tag(status, sub)
return VpgMembership(
vpg_identifier=str(pick(row, "VpgIdentifier", "vpgIdentifier") or ""),
vpg_name=str(pick(row, "VpgName", "vpgName") or ""),
status=status,
sub_status=sub,
can_tag=ok,
skip_reason=reason,
)
def group_vm_rows(rows: list[dict[str, Any]]) -> list[VmMatch]:
"""One GET /v1/vms row per VM-in-VPG. Group by VmIdentifier."""
by_id: dict[str, VmMatch] = {}
order: list[str] = []
for row in rows:
vm_id = pick(row, "VmIdentifier", "vmIdentifier")
vm_name = pick(row, "VmName", "vmName") or ""
if not vm_id:
continue
vm_id = str(vm_id)
if vm_id not in by_id:
by_id[vm_id] = VmMatch(vm_identifier=vm_id, vm_name=str(vm_name))
order.append(vm_id)
elif vm_name and not by_id[vm_id].vm_name:
by_id[vm_id].vm_name = str(vm_name)
vpg = membership_from_row(row)
if vpg.vpg_identifier and vpg.vpg_identifier not in {
x.vpg_identifier for x in by_id[vm_id].vpgs
}:
by_id[vm_id].vpgs.append(vpg)
return [by_id[i] for i in order]
def find_from_rows(query: str, rows: list[dict[str, Any]]) -> FindResult:
matches = group_vm_rows(rows)
if not matches:
return FindResult(
outcome="none",
query=query,
message=(
f"No protected VM matched {query!r}. "
"Unprotected or unknown: refuse the change. Zerto cannot rewind this."
),
)
if len(matches) > 1:
names = ", ".join(f"{m.vm_name} ({m.vm_identifier})" for m in matches)
return FindResult(
outcome="ambiguous",
query=query,
message=(
f"{len(matches)} VMs matched {query!r}: {names}. "
"Pass a Zerto vmIdentifier. Do not mutate."
),
matches=matches,
)
vm = matches[0]
if not vm.vpgs:
return FindResult(
outcome="none",
query=query,
message=(
f"VM {vm.vm_name} ({vm.vm_identifier}) has no VPG. "
"Unprotected: refuse the change."
),
vm=vm,
)
taggable = [v for v in vm.vpgs if v.can_tag]
skipped = [v for v in vm.vpgs if not v.can_tag]
skip_txt = ""
if skipped:
bits = "; ".join(f"{v.vpg_name}: {v.skip_reason}" for v in skipped)
skip_txt = f" Skipping: {bits}."
if not taggable:
return FindResult(
outcome="ok",
query=query,
message=(
f"VM {vm.vm_name} is in {len(vm.vpgs)} VPG(s) but none can be tagged right now."
f"{skip_txt} Refuse the change."
),
vm=vm,
)
return FindResult(
outcome="ok",
query=query,
message=(
f"VM {vm.vm_name} ({vm.vm_identifier}): "
f"{len(taggable)} protecting VPG(s) to tag, {len(skipped)} skipped."
f"{skip_txt}"
),
vm=vm,
)
+71
View File
@@ -0,0 +1,71 @@
"""FLR session helpers."""
from __future__ import annotations
import asyncio
from typing import Any
from zerto_rewind_mcp.client import ZertoClient, ZertoError
from zerto_rewind_mcp.util import pick
def session_id_from(payload: Any) -> str:
if payload is None:
raise ZertoError("FLR start returned empty body")
if isinstance(payload, str):
return payload.strip().strip('"')
if isinstance(payload, dict):
value = pick(
payload,
"sessionId",
"SessionId",
"flrSessionIdentifier",
"identifier",
"Identifier",
"id",
"Id",
)
if value:
return str(value)
raise ZertoError(f"Could not read FLR session id from {payload!r}")
def download_token_from(payload: Any) -> str:
if isinstance(payload, str):
return payload.strip().strip('"')
if isinstance(payload, dict):
value = pick(payload, "token", "Token", "downloadToken", "DownloadToken")
if value:
return str(value)
raise ZertoError(f"Could not read FLR download token from {payload!r}")
async def wait_flr_ready(
client: ZertoClient,
session_id: str,
*,
timeout_s: float = 180.0,
interval_s: float = 3.0,
) -> dict[str, Any]:
deadline = asyncio.get_event_loop().time() + timeout_s
last: Any = None
while asyncio.get_event_loop().time() < deadline:
last = await client.get_flr(session_id)
status = ""
if isinstance(last, dict):
status = str(pick(last, "Status", "status", "state", "State") or "")
if status.lower() in {"ready", "mounted", "available"}:
return last if isinstance(last, dict) else {"status": status}
if "mountinprogress" in status.lower() or "inprogress" in status.lower():
await asyncio.sleep(interval_s)
continue
if status.lower() in {"failed", "error"}:
raise ZertoError(f"FLR session {session_id} failed: {last}")
# some appliances omit status once mounted
if isinstance(last, dict) and not status:
return last
await asyncio.sleep(interval_s)
raise ZertoError(
f"FLR session {session_id} not ready within {timeout_s:.0f}s (last={last!r}). "
"Retry browse after a minute, or check EJC is not running and no clone/test is active."
)
+418
View File
@@ -0,0 +1,418 @@
"""Rewind-shaped MCP. One process for the demo. Not a second ZVM catalog."""
from __future__ import annotations
import json
from pathlib import Path
from typing import Any
from mcp.server.fastmcp import FastMCP
from zerto_rewind_mcp.catalog import CatalogEntry, MutatingCatalog
from zerto_rewind_mcp.checkpoints import checkpoint_id, checkpoint_tag, make_tag, tag_vpgs
from zerto_rewind_mcp.client import ZertoClient, ZertoError
from zerto_rewind_mcp.config import load_catalog, load_config
from zerto_rewind_mcp.protection import find_from_rows
from zerto_rewind_mcp.recover import download_token_from, session_id_from, wait_flr_ready
from zerto_rewind_mcp.util import pick
mcp = FastMCP("zerto_rewind_mcp")
_client: ZertoClient | None = None
_catalog: MutatingCatalog | None = None
_settings: dict[str, Any] = {}
def _dump(payload: Any) -> str:
return json.dumps(payload, indent=2, default=str)
def _ensure_settings() -> dict[str, Any]:
global _settings, _catalog
if not _settings:
_settings = load_config()
_catalog = load_catalog(_settings)
return _settings
def get_client() -> ZertoClient:
global _client
settings = _ensure_settings()
if _client is None:
if not settings.get("zerto_url") or not settings.get("username"):
raise ZertoError(
"Missing zerto_url/username. Copy config.example.json to config.json "
"or set ZERTO_URL and ZERTO_USERNAME."
)
_client = ZertoClient(
base_url=str(settings["zerto_url"]),
username=str(settings["username"]),
password=str(settings.get("password") or ""),
client_id=str(settings.get("client_id") or "zerto-client"),
verify_tls=bool(settings.get("verify_tls")),
)
return _client
def get_catalog() -> MutatingCatalog:
_ensure_settings()
assert _catalog is not None
return _catalog
async def _find(query: str):
client = get_client()
q = query.strip()
looks_id = len(q) >= 32 and "-" in q
rows = await client.get_vms(vm_identifier=q) if looks_id else await client.get_vms(vm_name=q)
if not rows and looks_id:
rows = await client.get_vms(vm_name=q)
elif not rows:
rows = await client.get_vms(vm_identifier=q)
return find_from_rows(q, rows)
@mcp.tool(
name="zerto_find_protection",
annotations={
"title": "Find Zerto protection for a VM",
"readOnlyHint": True,
"destructiveHint": False,
"idempotentHint": True,
"openWorldHint": True,
},
)
async def zerto_find_protection(query: str) -> str:
"""Resolve a VM name, hostname, or Zerto vmIdentifier to exactly one VM and every VPG it is in.
Zero matches or two-plus VMs: stop. Do not mutate. A VM can be in up to three VPGs.
"""
try:
result = await _find(query)
except ZertoError as exc:
return _dump({"ok": False, "message": str(exc)})
payload = result.as_dict()
payload["ok"] = result.outcome == "ok" and bool(result.taggable_vpgs)
return _dump(payload)
@mcp.tool(
name="zerto_create_tagged_checkpoint",
annotations={
"title": "Insert tagged checkpoints on every protecting VPG",
"readOnlyHint": False,
"destructiveHint": False,
"idempotentHint": False,
"openWorldHint": True,
},
)
async def zerto_create_tagged_checkpoint(
query: str,
change_id: str,
agent: str = "agent",
tag: str | None = None,
) -> str:
"""Insert the same tagged checkpoint on every protecting VPG, then wait.
Call this before every cataloged guest-mutating tool. If it fails, refuse the change.
Tag format defaults to ai:<agent>:<change_id>:<utc>.
Docs: tagged checkpoints are not supported when the protected site is Azure or AWS;
run this against the vSphere protected ZVM.
"""
try:
result = await _find(query)
if result.outcome != "ok" or not result.taggable_vpgs:
payload = result.as_dict()
payload["ok"] = False
return _dump(payload)
use_tag = tag or make_tag(agent, change_id)
out = await tag_vpgs(get_client(), result, use_tag)
return _dump(out)
except ZertoError as exc:
return _dump({"ok": False, "message": str(exc)})
@mcp.tool(
name="zerto_list_checkpoints",
annotations={
"title": "List journal checkpoints for a VPG",
"readOnlyHint": True,
"destructiveHint": False,
"idempotentHint": True,
"openWorldHint": True,
},
)
async def zerto_list_checkpoints(vpg_identifier: str, tag: str | None = None) -> str:
"""List checkpoints on a VPG. Pass tag to keep only matching tagged checkpoints."""
try:
rows = await get_client().list_checkpoints(vpg_identifier)
except ZertoError as exc:
return _dump({"ok": False, "message": str(exc)})
items = []
for row in rows:
item = {
"checkpoint_id": checkpoint_id(row),
"tag": checkpoint_tag(row),
"timestamp": pick(row, "Timestamp", "TimeStamp", "timestamp"),
}
if tag and item["tag"] != tag:
continue
items.append(item)
return _dump(
{
"ok": True,
"vpg_identifier": vpg_identifier,
"count": len(items),
"checkpoints": items,
}
)
@mcp.tool(
name="zerto_guard_before_mutate",
annotations={
"title": "Guard: find protection and tag before a mutating call",
"readOnlyHint": False,
"destructiveHint": False,
"idempotentHint": False,
"openWorldHint": True,
},
)
async def zerto_guard_before_mutate(
query: str,
change_id: str,
agent: str = "agent",
) -> str:
"""Find protection and tag. If ok is not true, do not mutate.
The skill and any cataloged tool wrapper call this first. Capture-before-execute:
if this does not return ok=true, do not run the guest-mutating call.
"""
return await zerto_create_tagged_checkpoint(query=query, change_id=change_id, agent=agent)
@mcp.tool(
name="zerto_list_mutating_tools",
annotations={
"title": "List the mutating catalog",
"readOnlyHint": True,
"destructiveHint": False,
"idempotentHint": True,
"openWorldHint": False,
},
)
async def zerto_list_mutating_tools() -> str:
"""Tools that must be guarded. Unlisted tools pass through."""
entries = [e.as_dict() for e in get_catalog().list()]
return _dump({"ok": True, "count": len(entries), "mutating_tools": entries})
@mcp.tool(
name="zerto_add_mutating_tool",
annotations={
"title": "Add a tool to the mutating catalog",
"readOnlyHint": False,
"destructiveHint": False,
"idempotentHint": True,
"openWorldHint": False,
},
)
async def zerto_add_mutating_tool(
server: str,
tool: str,
vm_arg: str,
notes: str = "",
) -> str:
"""Register an MCP tool as guest-mutating. vm_arg holds VM name/hostname/id."""
try:
entry = CatalogEntry(server=server, tool=tool, vm_arg=vm_arg, notes=notes)
get_catalog().add(entry)
except ValueError as exc:
return _dump({"ok": False, "message": str(exc)})
return _dump({"ok": True, "entry": entry.as_dict()})
@mcp.tool(
name="zerto_recover_file",
annotations={
"title": "FLR: pull a file from a tagged checkpoint",
"readOnlyHint": False,
"destructiveHint": True,
"idempotentHint": False,
"openWorldHint": True,
},
)
async def zerto_recover_file(
vpg_identifier: str,
vm_identifier: str,
checkpoint_identifier: str,
guest_path: str,
confirmed: bool = False,
dest_dir: str | None = None,
) -> str:
"""File-level recovery from a journal checkpoint. The VM stays up.
Requires confirmed=true (human yes). Cannot run during clone/test/live/EJC.
10.9 FLR Operator role fails; use an Administrator account.
"""
if not confirmed:
return _dump(
{
"ok": False,
"needs_confirm": True,
"message": (
"Set confirmed=true after a human yes. "
"FLR mounts a disk on the recovery site."
),
}
)
client = get_client()
dest = Path(dest_dir or _settings.get("recovery_dir") or "./recovered")
dest.mkdir(parents=True, exist_ok=True)
session_id = None
try:
started = await client.start_flr(
vpg_identifier,
vm_identifier,
checkpoint_identifier,
str(dest.resolve()),
)
session_id = session_id_from(started)
await wait_flr_ready(client, session_id)
await client.browse_flr(session_id, path="", recursive=False)
token_payload = await client.download_flr(session_id, [guest_path])
token = download_token_from(token_payload)
blob = await client.fetch_download(token)
name = Path(guest_path.replace("\\", "/")).name or "recovered.bin"
out_path = dest / name
out_path.write_bytes(blob)
return _dump(
{
"ok": True,
"path": str(out_path.resolve()),
"bytes": len(blob),
"session_id": session_id,
"message": f"Wrote {len(blob)} bytes to {out_path}",
}
)
except ZertoError as exc:
return _dump({"ok": False, "message": str(exc)})
finally:
if session_id:
try:
await client.end_flr(session_id)
except ZertoError:
pass
@mcp.tool(
name="zerto_start_failover_test",
annotations={
"title": "Start a failover test (inspect only)",
"readOnlyHint": False,
"destructiveHint": True,
"idempotentHint": False,
"openWorldHint": True,
},
)
async def zerto_start_failover_test(
vpg_identifier: str,
confirmed: bool = False,
checkpoint_identifier: str | None = None,
) -> str:
"""Spin up test VMs from a checkpoint. Not a production failback. Requires confirmed=true."""
if not confirmed:
return _dump(
{
"ok": False,
"needs_confirm": True,
"message": (
"Set confirmed=true after a human yes. "
"Failover test creates VMs on the test network."
),
}
)
try:
result = await get_client().start_failover_test(vpg_identifier, checkpoint_identifier)
return _dump({"ok": True, "result": result})
except ZertoError as exc:
return _dump({"ok": False, "message": str(exc)})
@mcp.tool(
name="zerto_stop_failover_test",
annotations={
"title": "Stop a failover test",
"readOnlyHint": False,
"destructiveHint": True,
"idempotentHint": False,
"openWorldHint": True,
},
)
async def zerto_stop_failover_test(
vpg_identifier: str,
success: bool = True,
summary: str = "stopped by zerto rewind mcp",
confirmed: bool = False,
) -> str:
"""Stop an in-progress failover test. Requires confirmed=true."""
if not confirmed:
return _dump(
{
"ok": False,
"needs_confirm": True,
"message": "Set confirmed=true after a human yes.",
}
)
try:
result = await get_client().stop_failover_test(vpg_identifier, success, summary)
return _dump({"ok": True, "result": result})
except ZertoError as exc:
return _dump({"ok": False, "message": str(exc)})
@mcp.tool(
name="zerto_offsite_clone",
annotations={
"title": "Offsite clone from a checkpoint",
"readOnlyHint": False,
"destructiveHint": True,
"idempotentHint": False,
"openWorldHint": True,
},
)
async def zerto_offsite_clone(
vpg_identifier: str,
checkpoint_identifier: str,
confirmed: bool = False,
datastore_identifier: str | None = None,
) -> str:
"""Copy VMs at the recovery site from a checkpoint. Prod stays up. Requires confirmed=true."""
if not confirmed:
return _dump(
{
"ok": False,
"needs_confirm": True,
"message": (
"Set confirmed=true after a human yes. "
"Clone consumes recovery datastore."
),
}
)
try:
result = await get_client().start_clone(
vpg_identifier,
checkpoint_identifier,
datastore_identifier=datastore_identifier,
)
return _dump({"ok": True, "result": result})
except ZertoError as exc:
return _dump({"ok": False, "message": str(exc)})
def main() -> None:
mcp.run()
if __name__ == "__main__":
main()
+133
View File
@@ -0,0 +1,133 @@
"""VPG status / substatus integers from GET /v1/vpgs/statuses and /substatuses.
Live 10.9 ZVM (192.168.50.30, 2026-09-21). Older api_lessons mapping
(0=Protecting, 1=Moving) is wrong on this appliance.
"""
# GET /v1/vpgs/statuses
INITIALIZING = 0
MEETING_SLA = 1
NOT_MEETING_SLA = 2
RPO_NOT_MEETING_SLA = 3
HISTORY_NOT_MEETING_SLA = 4
FAILING_OVER = 5
MOVING = 6
DELETING = 7
RECOVERED = 8
STATUS_NAME = {
INITIALIZING: "Initializing",
MEETING_SLA: "MeetingSLA",
NOT_MEETING_SLA: "NotMeetingSLA",
RPO_NOT_MEETING_SLA: "RpoNotMeetingSLA",
HISTORY_NOT_MEETING_SLA: "HistoryNotMeetingSLA",
FAILING_OVER: "FailingOver",
MOVING: "Moving",
DELETING: "Deleting",
RECOVERED: "Recovered",
}
# GET /v1/vpgs/substatuses
NONE = 0
INITIAL_SYNC = 1
CREATING = 2
VOLUME_INITIAL_SYNC = 3
SYNC = 4
RECOVERY_POSSIBLE = 5
DELTA_SYNC = 6
NEEDS_CONFIGURATION = 7
ERROR = 8
EMPTY_PROTECTION_GROUP = 9
DISCONNECTED_NO_RECOVERY_POINTS = 10
FULL_SYNC = 11
VOLUME_DELTA_SYNC = 12
VOLUME_FULL_SYNC = 13
FAILING_OVER_COMMITTING = 14
FAILING_OVER_BEFORE_COMMIT = 15
FAILING_OVER_ROLLING_BACK = 16
PROMOTING = 17
MOVING_COMMITTING = 18
MOVING_BEFORE_COMMIT = 19
MOVING_ROLLING_BACK = 20
DELETING_SUB = 21
PENDING_REMOVE = 22
BITMAP_SYNC = 23
DISCONNECTED_FROM_PEER = 24
REPLICATION_PAUSED_USER = 25
REPLICATION_PAUSED_SYSTEM = 26
ADDED_VMS_IN_INITIAL_SYNC = 34
SUBSTATUS_NAME = {
NONE: "None",
INITIAL_SYNC: "InitialSync",
CREATING: "Creating",
VOLUME_INITIAL_SYNC: "VolumeInitialSync",
SYNC: "Sync",
RECOVERY_POSSIBLE: "RecoveryPossible",
DELTA_SYNC: "DeltaSync",
NEEDS_CONFIGURATION: "NeedsConfiguration",
ERROR: "Error",
EMPTY_PROTECTION_GROUP: "EmptyProtectionGroup",
DISCONNECTED_NO_RECOVERY_POINTS: "DisconnectedFromPeerNoRecoveryPoints",
FULL_SYNC: "FullSync",
VOLUME_DELTA_SYNC: "VolumeDeltaSync",
VOLUME_FULL_SYNC: "VolumeFullSync",
FAILING_OVER_COMMITTING: "FailingOverCommitting",
FAILING_OVER_BEFORE_COMMIT: "FailingOverBeforeCommit",
FAILING_OVER_ROLLING_BACK: "FailingOverRollingBack",
PROMOTING: "Promoting",
MOVING_COMMITTING: "MovingCommitting",
MOVING_BEFORE_COMMIT: "MovingBeforeCommit",
MOVING_ROLLING_BACK: "MovingRollingBack",
DELETING_SUB: "Deleting",
PENDING_REMOVE: "PendingRemove",
BITMAP_SYNC: "BitmapSync",
DISCONNECTED_FROM_PEER: "DisconnectedFromPeer",
REPLICATION_PAUSED_USER: "ReplicationPausedUserInitiated",
REPLICATION_PAUSED_SYSTEM: "ReplicationPausedSystemInitiated",
ADDED_VMS_IN_INITIAL_SYNC: "AddedVmsInInitialSync",
}
SYNCING = {
INITIAL_SYNC,
CREATING,
VOLUME_INITIAL_SYNC,
SYNC,
DELTA_SYNC,
FULL_SYNC,
VOLUME_DELTA_SYNC,
VOLUME_FULL_SYNC,
BITMAP_SYNC,
ADDED_VMS_IN_INITIAL_SYNC,
}
TAGGABLE_STATUS = {
MEETING_SLA,
NOT_MEETING_SLA,
RPO_NOT_MEETING_SLA,
HISTORY_NOT_MEETING_SLA,
}
def status_name(value: int | None) -> str:
if value is None:
return "unknown"
return STATUS_NAME.get(value, f"status:{value}")
def substatus_name(value: int | None) -> str:
if value is None:
return "unknown"
return SUBSTATUS_NAME.get(value, f"substatus:{value}")
def can_tag(status: int | None, substatus: int | None) -> tuple[bool, str | None]:
"""Whether a tagged checkpoint can be inserted on this VPG right now."""
if status not in TAGGABLE_STATUS:
return False, f"VPG status is {status_name(status)}; not MeetingSLA"
if substatus in SYNCING:
return False, (
f"VPG is {substatus_name(substatus)}; "
"checkpoints are not durable until sync ends"
)
return True, None
+16
View File
@@ -0,0 +1,16 @@
"""JSON key access. Zerto mixes PascalCase and camelCase across versions."""
from __future__ import annotations
from typing import Any
def pick(row: dict[str, Any], *names: str) -> Any:
for name in names:
if name in row and row[name] is not None:
return row[name]
lower = {k.lower(): v for k, v in row.items()}
for name in names:
if name.lower() in lower and lower[name.lower()] is not None:
return lower[name.lower()]
return None