commit 38ba9c1b5097ae87a7ace44f110c3f31302c0836 Author: claude Date: Mon Sep 21 12:09:11 2026 -0400 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. diff --git a/.claude/gitea-ship.json b/.claude/gitea-ship.json new file mode 100644 index 0000000..53a70cb --- /dev/null +++ b/.claude/gitea-ship.json @@ -0,0 +1,10 @@ +{ + "owner": "justin", + "registry_lan": "192.168.0.2:1234", + "registry_fqdn": "git.jpaul.io", + "runner_label": "docker", + "deploy": "ssh-systemd", + "version_source": "semver-code", + "migrations": "none", + "notes": "PoC MCP. Not containerized. Lab ZVM credentials never go in git (config.json is gitignored)." +} diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..9421fda --- /dev/null +++ b/.gitignore @@ -0,0 +1,17 @@ +.venv/ +__pycache__/ +*.py[cod] +*.egg-info/ +dist/ +build/ +.pytest_cache/ +.ruff_cache/ +.coverage +htmlcov/ +.env +config.json +catalog.json +*.log +.mcp.json +.claude/* +!.claude/gitea-ship.json diff --git a/CONTEXT.md b/CONTEXT.md new file mode 100644 index 0000000..dc5c3cc --- /dev/null +++ b/CONTEXT.md @@ -0,0 +1,41 @@ +# Zerto AI Rewind + +PoC MCP that teaches an agent to discover Zerto protection, pin a tagged checkpoint before changing a VM, and recover a file from that tag. If the loop works, these tools are the delta to put in official ZVM MCP. + +## Language + +**Tagged checkpoint**: +A named bookmark in a VPG journal, inserted by `POST /v1/vpgs/{id}/checkpoints`. Crash-consistent write-order only; not application-quiesced unless someone scripted that separately. +_Avoid_: user checkpoint, snapshot, backup, restore point (unqualified) + +**VPG**: +A Virtual Protection Group. One to many VMs sharing a journal. A VM can belong to at most three VPGs, recovered to different sites. +_Avoid_: job, policy, replication group + +**Rewind**: +The agent loop: find protection, tag every protecting VPG, mutate, then bounded recover. Not a Zerto product name. +_Avoid_: failover (that's DR), undo (that's git or Moholo) + +**Bounded recover**: +FLR, offsite clone, or failover test. Failover Live is not a rewind tool. +_Avoid_: recover (unqualified), fail back, restore the VPG + +**File-level recovery (FLR)**: +Mount a VM from a journal checkpoint and pull files. The VM stays up. 10.9 FLR Operator RBAC is broken; Administrator is the documented workaround. +_Avoid_: file restore (unqualified), instant restore (local-journal VMs only, not v1) + +**find_protection**: +Resolve a VM name, hostname, or Zerto vmIdentifier to exactly one VM and every VPG it is in. Zero or two-plus VMs is a hard stop. +_Avoid_: GetVms (that's the raw inventory call) + +**Protecting VPG**: +A VPG whose status is MeetingSLA or a NotMeetingSLA variant, and whose substatus is not a sync. Only these get tagged. 10.9 status 0 is Initializing, not Protecting. A resync deletes existing checkpoints. +_Avoid_: healthy, in sync, Protecting (as status 0) + +**Mutating catalog**: +The opt-in list of MCP tools that must call `zerto_guard_before_mutate` first. Unlisted tools pass through. Users add entries; the starter list is not the whole world. +_Avoid_: denylist, hold-everything + +**Official ZVM MCP**: +HPE Zerto 10.9 MCP (`ZVM.MCP`): inventory, VPG settings, failover test. Not in the demo path. This PoC is one server. +_Avoid_: Zerto MCP (unqualified when you mean this repo) diff --git a/README.md b/README.md new file mode 100644 index 0000000..05235e5 --- /dev/null +++ b/README.md @@ -0,0 +1,74 @@ +# zerto-ai-rewind + +PoC MCP that makes an agent pin a Zerto tagged checkpoint before it changes a VM, then pull a file back from that tag. + +If the loop works, these tools are the delta to add to official ZVM MCP (`ZVM.MCP`, 10.9). This repo is one MCP process for the demo. It is not a second full ZVM catalog. + +## What it does + +1. `zerto_find_protection` — VM name, hostname, or vmIdentifier to exactly one VM and every VPG. Zero or two-plus VMs: stop. +2. `zerto_create_tagged_checkpoint` / `zerto_guard_before_mutate` — same tag on every protecting VPG, wait until listed. +3. `zerto_recover_file` — FLR after a human sets `confirmed=true`. +4. Mutating catalog — opt-in list of MCP tools that must be guarded. Unlisted tools pass through. Users add entries. + +Official ZVM MCP already has inventory and failover test. It does not insert tagged checkpoints or run FLR. + +## Setup + +Python 3.12+. + +```bash +python3 -m venv .venv +source .venv/bin/activate +pip install -e ".[dev]" +cp config.example.json config.json +# edit zerto_url, username, password +``` + +Keycloak password-grant, client_id `zerto-client` on 10.x. Appliance certs are self-signed; `verify_tls` defaults to false. + +stdio MCP (Claude Desktop, VS Code, Cursor, OpenCode): + +```json +{ + "mcpServers": { + "zerto-rewind": { + "command": "zerto-rewind-mcp", + "env": { + "ZERTO_REWIND_CONFIG": "/absolute/path/to/config.json" + } + } + } +} +``` + +Copy `skills/zerto-rewind/SKILL.md` into the client's skill path. + +```bash +pytest +``` + +## Demo + +Protected app VM. Agent is about to edit a guest config file. + +1. Guard: discover VPG set, insert tagged checkpoint, wait. +2. Agent writes the bad config. +3. Human confirms. +4. `zerto_recover_file` from that tag. + +Git never had the file. RPO is the journal, not last night's backup. + +## Certified for this PoC + +vSphere ZVM 10.x and ZCA on AWS/Azure, same REST paths. HVM is out (separate swagger). Failover Live is not a tool. + +Tagged checkpoints cannot be inserted when the **protected** site is Azure or AWS (Zerto API). Point this server at the vSphere protected ZVM. + +## Not this product + +[Moholo Agent Rewind](https://github.com/moholo-founder/agent-rewind) snapshots the agent's laptop tools. This server uses the Zerto journal as the snapshot store. Do not copy their file blobs. + +## Upstream + +Ask Zerto engineering to add to `ZVM.MCP`: find-by-unique-VM-with-all-VPGs, tagged checkpoint insert that waits, FLR. Keep VPG settings CRUD where it already is. diff --git a/config.example.json b/config.example.json new file mode 100644 index 0000000..1fe7dac --- /dev/null +++ b/config.example.json @@ -0,0 +1,22 @@ +{ + "zerto_url": "https://zvm.example.com", + "username": "rewind-svc", + "password": "change-me", + "client_id": "zerto-client", + "verify_tls": false, + "recovery_dir": "./recovered", + "mutating_tools": [ + { + "server": "ssh", + "tool": "exec", + "vm_arg": "host", + "notes": "Remote shell on a guest. vm_arg is the hostname." + }, + { + "server": "ansible", + "tool": "run_playbook", + "vm_arg": "limit", + "notes": "Playbook target host/group. Resolve to a single VM before guard." + } + ] +} diff --git a/docs/recover-ladder.md b/docs/recover-ladder.md new file mode 100644 index 0000000..cfa2fda --- /dev/null +++ b/docs/recover-ladder.md @@ -0,0 +1,57 @@ +# Recover ladder: FLR vs whole-VM rewind + +Zerto's journal can rewind a file or a whole VM. The agent picks the smallest +operation that actually undoes the damage. Failover Live is DR, not the default +undo. + +## File-level recovery (FLR) + +Use when the guest still boots, SSH/WinRM still works, and the damage is a +known path (config, dropped file, one directory). + +FLR mounts a checkpoint and copies files out. The protected VM stays up. +Official API: `POST /v1/flrs` then browse/download. This MCP writes the file +to `recovery_dir` on the MCP host. Putting it back on the guest is a second +step (scp/ssh). That copy-back is not Zerto; it is ordinary file transfer. + +Do not use FLR when: + +- `EnabledActions.IsFlrEnabled` is false (initial sync, clone, test, live, + move, or EJC running) +- The guest cannot boot or accept a file +- You do not know which files changed (package install, kernel, ransomware) +- Linux file >1.5GB, or the name has `\ / : * ? " < > |` +- OS-level dedup volumes +- 10.9 FLR Operator role (broken; Administrator is the documented workaround) + +A tagged checkpoint must already exist. Initial sync has an empty journal +(`GET .../checkpoints` returns `[]`). Guard refuses until status is MeetingSLA +(or NotMeetingSLA) and substatus is not a sync. + +## Whole-VM rewind + +Use when FLR cannot put the guest back: OS broken, too many files, services or +packages, unknown blast radius. + +| Operation | What it does | When | +|---|---|---| +| Failover test | Test VMs from a checkpoint. Prod stays up. | Inspect only | +| Offsite clone | Copy VMs at recovery, powered off, unprotected | Inspect or graft | +| Instant restore | Local-journal VPGs only; one VM, journal kept | Same-site local VPG | +| Failover Live | Production DR. Reverse protection. Stops the protected VM | Site or VM beyond clone/FLR, human confirmed | + +Failover Live is not the coding-agent oops button. Same-site VPGs still run a +real failover. Confirm in the client before the tool runs. + +## Guard vs recover + +1. `zerto_find_protection` must return exactly one VM and at least one + taggable VPG. +2. `zerto_guard_before_mutate` inserts the tagged checkpoint and waits until + it is listed. If that fails, do not change the guest. +3. After a bad change: FLR if the path is known and the guest is up; clone or + test to inspect; Failover Live only with a human yes. + +Status integers on 10.9 (`GET /v1/vpgs/statuses`): 0 Initializing, 1 +MeetingSLA, 2 NotMeetingSLA. Older notes that said 0=Protecting / 1=Moving are +wrong on this API. diff --git a/mcp.example.json b/mcp.example.json new file mode 100644 index 0000000..2322bdf --- /dev/null +++ b/mcp.example.json @@ -0,0 +1,10 @@ +{ + "mcpServers": { + "zerto-rewind": { + "command": "zerto-rewind-mcp", + "env": { + "ZERTO_REWIND_CONFIG": "/absolute/path/to/config.json" + } + } + } +} diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..d1b5ee5 --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,43 @@ +[project] +name = "zerto-rewind-mcp" +version = "0.1.0" +description = "PoC MCP: find Zerto protection, tagged checkpoints before guest mutations, FLR from that tag" +readme = "README.md" +requires-python = ">=3.12" +dependencies = [ + "httpx>=0.27", + "mcp[cli]>=1.9,<2", + "pydantic>=2.7", + "pydantic-settings>=2.4", +] + +[project.optional-dependencies] +dev = [ + "pytest>=8.2", + "pytest-asyncio>=0.23", + "ruff>=0.6", +] + +[project.scripts] +zerto-rewind-mcp = "zerto_rewind_mcp.server:main" + +[build-system] +requires = ["setuptools>=69"] +build-backend = "setuptools.build_meta" + +[tool.setuptools.packages.find] +where = ["src"] + +[tool.ruff] +line-length = 100 +target-version = "py312" +src = ["src", "tests"] + +[tool.ruff.lint] +extend-select = ["E", "F", "I", "B", "UP", "ASYNC", "RUF"] +ignore = ["ASYNC240"] + +[tool.pytest.ini_options] +asyncio_mode = "auto" +testpaths = ["tests"] +pythonpath = ["src"] diff --git a/skills/zerto-rewind/SKILL.md b/skills/zerto-rewind/SKILL.md new file mode 100644 index 0000000..835a544 --- /dev/null +++ b/skills/zerto-rewind/SKILL.md @@ -0,0 +1,54 @@ +--- +name: zerto-rewind +description: Before changing a VM, find its Zerto VPGs and insert a tagged checkpoint. Recover files from that tag with FLR after a human confirms. Use whenever an agent will mutate a guest that might be protected by Zerto. +--- + +# Zerto rewind + +Zerto already journals the VM. This skill makes the agent use that journal. Git does not have the guest file. Official ZVM MCP does not insert tagged checkpoints. + +You talk to **one** MCP: `zerto_rewind_mcp`. Do not also require official ZVM MCP. + +## Loop (mandatory) + +Before **every** guest-mutating tool call: + +1. Take the hostname / VM name / Zerto `vmIdentifier` from the tool args. +2. Call `zerto_guard_before_mutate` (or `zerto_find_protection` then `zerto_create_tagged_checkpoint`). +3. If `ok` is not true: **stop**. Do not mutate. +4. Then run the mutating call. + +Reads skip the guard. + +Unlisted MCP tools pass through. If you are about to change a protected VM with a tool that is not in the catalog, call `zerto_add_mutating_tool` (server, tool, `vm_arg`) and then guard. + +## find_protection outcomes + +| outcome | what you do | +|---|---| +| none | Unprotected or unknown. Refuse the change. Say Zerto cannot rewind this. | +| ambiguous | Two or more VMs matched. Ask for a `vmIdentifier`. Do not guess. | +| ok, no taggable VPG | Syncing or not Protecting. Refuse. A resync deletes checkpoints. | +| ok, taggable VPGs | Tag **every** protecting VPG with the same tag. Wait until listed (the tool blocks). | + +A VM can be in up to three VPGs (local backup + remote DR is common). Tag all of them. + +## Recover + +Human must confirm. Pass `confirmed=true` only after they say yes. + +- Bad config / dropped file: `zerto_recover_file` from **that tag**. +- Inspect a whole VM: `zerto_offsite_clone` or `zerto_start_failover_test`. +- Never Failover Live. Never Move. Those are DR, not rewind. + +## Facts that bite + +- A tagged checkpoint is crash-consistent, not app-quiesced. +- Tagged checkpoints are not supported when the **protected** site is Azure or AWS. Talk to the vSphere protected ZVM. +- 10.9 FLR Operator RBAC fails; Administrator is the documented workaround. +- FLR cannot run during clone, test, live failover, or EJC. +- Linux FLR: files >1.5GB are a bad idea; some characters in names are refused. + +## Tag + +Default: `ai:::`. Same string on every VPG for that call. diff --git a/src/zerto_rewind_mcp/__init__.py b/src/zerto_rewind_mcp/__init__.py new file mode 100644 index 0000000..fd9fb2a --- /dev/null +++ b/src/zerto_rewind_mcp/__init__.py @@ -0,0 +1,3 @@ +"""Zerto AI Rewind PoC MCP.""" + +__version__ = "0.1.0" diff --git a/src/zerto_rewind_mcp/__main__.py b/src/zerto_rewind_mcp/__main__.py new file mode 100644 index 0000000..8d00bb3 --- /dev/null +++ b/src/zerto_rewind_mcp/__main__.py @@ -0,0 +1,4 @@ +from zerto_rewind_mcp.server import main + +if __name__ == "__main__": + main() diff --git a/src/zerto_rewind_mcp/catalog.py b/src/zerto_rewind_mcp/catalog.py new file mode 100644 index 0000000..7a3123b --- /dev/null +++ b/src/zerto_rewind_mcp/catalog.py @@ -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) diff --git a/src/zerto_rewind_mcp/checkpoints.py b/src/zerto_rewind_mcp/checkpoints.py new file mode 100644 index 0000000..bf492b4 --- /dev/null +++ b/src/zerto_rewind_mcp/checkpoints.py @@ -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, + } diff --git a/src/zerto_rewind_mcp/client.py b/src/zerto_rewind_mcp/client.py new file mode 100644 index 0000000..3d73849 --- /dev/null +++ b/src/zerto_rewind_mcp/client.py @@ -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, + ) diff --git a/src/zerto_rewind_mcp/config.py b/src/zerto_rewind_mcp/config.py new file mode 100644 index 0000000..5ae1064 --- /dev/null +++ b/src/zerto_rewind_mcp/config.py @@ -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) diff --git a/src/zerto_rewind_mcp/protection.py b/src/zerto_rewind_mcp/protection.py new file mode 100644 index 0000000..0e50d2a --- /dev/null +++ b/src/zerto_rewind_mcp/protection.py @@ -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, + ) diff --git a/src/zerto_rewind_mcp/recover.py b/src/zerto_rewind_mcp/recover.py new file mode 100644 index 0000000..be80b84 --- /dev/null +++ b/src/zerto_rewind_mcp/recover.py @@ -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." + ) diff --git a/src/zerto_rewind_mcp/server.py b/src/zerto_rewind_mcp/server.py new file mode 100644 index 0000000..1a2bec0 --- /dev/null +++ b/src/zerto_rewind_mcp/server.py @@ -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:::. + 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() diff --git a/src/zerto_rewind_mcp/status.py b/src/zerto_rewind_mcp/status.py new file mode 100644 index 0000000..ff25f6a --- /dev/null +++ b/src/zerto_rewind_mcp/status.py @@ -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 diff --git a/src/zerto_rewind_mcp/util.py b/src/zerto_rewind_mcp/util.py new file mode 100644 index 0000000..0325075 --- /dev/null +++ b/src/zerto_rewind_mcp/util.py @@ -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 diff --git a/tests/test_catalog.py b/tests/test_catalog.py new file mode 100644 index 0000000..2675394 --- /dev/null +++ b/tests/test_catalog.py @@ -0,0 +1,33 @@ +import json +from pathlib import Path + +from zerto_rewind_mcp.catalog import CatalogEntry, MutatingCatalog, entry_from_dict + + +def test_add_and_get(tmp_path: Path): + path = tmp_path / "config.json" + path.write_text( + json.dumps({"zerto_url": "https://zvm", "mutating_tools": []}), + encoding="utf-8", + ) + cat = MutatingCatalog(path=path) + cat.add(CatalogEntry(server="ssh", tool="exec", vm_arg="host", notes="guest shell")) + found = cat.get("SSH", "Exec") + assert found is not None + assert found.vm_arg == "host" + saved = json.loads(path.read_text(encoding="utf-8")) + assert saved["zerto_url"] == "https://zvm" + assert saved["mutating_tools"][0]["tool"] == "exec" + + +def test_unlisted_is_missing(): + cat = MutatingCatalog() + assert cat.get("foo", "bar") is None + + +def test_entry_requires_fields(): + try: + entry_from_dict({"server": "ssh"}) + raise AssertionError("expected ValueError") + except ValueError: + pass diff --git a/tests/test_checkpoints.py b/tests/test_checkpoints.py new file mode 100644 index 0000000..4eb5e2e --- /dev/null +++ b/tests/test_checkpoints.py @@ -0,0 +1,24 @@ +from datetime import UTC, datetime + +from zerto_rewind_mcp.checkpoints import checkpoint_id, checkpoint_tag, make_tag + + +def test_tag_format(): + when = datetime(2026, 9, 21, 15, 4, 5, tzinfo=UTC) + tag = make_tag("codex", "chg-99", when) + assert tag == "ai:codex:chg-99:20260921T150405Z" + + +def test_tag_strips_junk(): + tag = make_tag("agent one", "path/with spaces", datetime(2026, 1, 1, tzinfo=UTC)) + assert " " not in tag + assert tag.startswith("ai:agent-one:path-with-spaces:") + + +def test_checkpoint_row_keys(): + row = {"CheckpointId": "cp-1", "Tag": "ai:x:y:z", "Timestamp": 1} + assert checkpoint_id(row) == "cp-1" + assert checkpoint_tag(row) == "ai:x:y:z" + row2 = {"checkpointId": "cp-2", "tag": "t"} + assert checkpoint_id(row2) == "cp-2" + assert checkpoint_tag(row2) == "t" diff --git a/tests/test_protection.py b/tests/test_protection.py new file mode 100644 index 0000000..56f7525 --- /dev/null +++ b/tests/test_protection.py @@ -0,0 +1,103 @@ +from zerto_rewind_mcp.protection import find_from_rows, group_vm_rows +from zerto_rewind_mcp.status import ( + BITMAP_SYNC, + DELTA_SYNC, + INITIAL_SYNC, + INITIALIZING, + MEETING_SLA, + NOT_MEETING_SLA, +) + + +def _row(vm_id, vm_name, vpg_id, vpg_name, status=MEETING_SLA, sub=0): + return { + "VmIdentifier": vm_id, + "VmName": vm_name, + "VpgIdentifier": vpg_id, + "VpgName": vpg_name, + "Status": status, + "SubStatus": sub, + } + + +def test_none(): + result = find_from_rows("web01", []) + assert result.outcome == "none" + assert "refuse" in result.message.lower() + + +def test_unique_two_vpgs(): + rows = [ + _row("vm-1", "web01", "vpg-local", "local-bu"), + _row("vm-1", "web01", "vpg-dr", "dr-remote"), + ] + result = find_from_rows("web01", rows) + assert result.outcome == "ok" + assert result.vm is not None + assert result.vm.vm_identifier == "vm-1" + assert len(result.vm.vpgs) == 2 + assert len(result.taggable_vpgs) == 2 + + +def test_ambiguous_two_vms(): + rows = [ + _row("vm-1", "web01", "vpg-a", "a"), + _row("vm-2", "web01", "vpg-b", "b"), + ] + result = find_from_rows("web01", rows) + assert result.outcome == "ambiguous" + assert len(result.matches) == 2 + assert result.taggable_vpgs == [] + + +def test_skip_syncing_vpg(): + rows = [ + _row("vm-1", "app", "vpg-ok", "ok", status=MEETING_SLA, sub=0), + _row("vm-1", "app", "vpg-sync", "syncing", status=MEETING_SLA, sub=INITIAL_SYNC), + _row("vm-1", "app", "vpg-delta", "delta", status=MEETING_SLA, sub=DELTA_SYNC), + _row("vm-1", "app", "vpg-bitmap", "bitmap", status=MEETING_SLA, sub=BITMAP_SYNC), + ] + result = find_from_rows("app", rows) + assert result.outcome == "ok" + names = {v.vpg_name for v in result.taggable_vpgs} + assert names == {"ok"} + assert "InitialSync" in result.message + + +def test_not_meeting_sla_is_taggable(): + rows = [_row("vm-1", "app", "vpg-a", "a", status=NOT_MEETING_SLA, sub=0)] + result = find_from_rows("app", rows) + assert result.outcome == "ok" + assert len(result.taggable_vpgs) == 1 + + +def test_initializing_not_taggable(): + rows = [_row("vm-1", "app", "vpg-a", "a", status=INITIALIZING, sub=INITIAL_SYNC)] + result = find_from_rows("app", rows) + assert result.outcome == "ok" + assert result.taggable_vpgs == [] + assert "refuse" in result.message.lower() + + +def test_all_syncing_ok_but_no_taggable(): + rows = [_row("vm-1", "app", "vpg-a", "a", status=MEETING_SLA, sub=INITIAL_SYNC)] + result = find_from_rows("app", rows) + assert result.outcome == "ok" + assert result.taggable_vpgs == [] + assert "refuse" in result.message.lower() + + +def test_camel_case_keys(): + rows = [ + { + "vmIdentifier": "id-9", + "vmName": "db01", + "vpgIdentifier": "v1", + "vpgName": "db-vpg", + "status": 1, + "subStatus": 0, + } + ] + grouped = group_vm_rows(rows) + assert grouped[0].vm_name == "db01" + assert grouped[0].vpgs[0].can_tag diff --git a/tests/test_recover.py b/tests/test_recover.py new file mode 100644 index 0000000..c2d52cc --- /dev/null +++ b/tests/test_recover.py @@ -0,0 +1,12 @@ +from zerto_rewind_mcp.recover import download_token_from, session_id_from + + +def test_session_id_shapes(): + assert session_id_from("abc") == "abc" + assert session_id_from({"sessionId": "s1"}) == "s1" + assert session_id_from({"Identifier": "s2"}) == "s2" + + +def test_download_token_shapes(): + assert download_token_from("tok") == "tok" + assert download_token_from({"downloadToken": "t2"}) == "t2"