571 lines
18 KiB
Python
571 lines
18 KiB
Python
"""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,
|
|
is_live_session,
|
|
resolve_flr_path,
|
|
session_id_from,
|
|
session_id_of,
|
|
session_rows,
|
|
session_summary,
|
|
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",
|
|
action: str = "",
|
|
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.
|
|
|
|
action: say plainly what you are about to do to this VM, e.g.
|
|
"edit /etc/nginx/nginx.conf" or "apt upgrade". It is written into the
|
|
checkpoint name so an operator reading the journal in the Zerto UI can see
|
|
which agent made the checkpoint and why. CheckpointName is the only field
|
|
the Zerto API accepts, so this is the only place that context can live.
|
|
|
|
Name format: ai:<agent> | <action> | vm=<vm> | change=<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,
|
|
action=action,
|
|
vm_name=result.vm.vm_name if result.vm else None,
|
|
)
|
|
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",
|
|
action: str = "",
|
|
) -> 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.
|
|
|
|
Pass action describing the change you are about to make; it is recorded in
|
|
the checkpoint name so the journal says which agent did what.
|
|
"""
|
|
return await zerto_create_tagged_checkpoint(
|
|
query=query, change_id=change_id, agent=agent, action=action
|
|
)
|
|
|
|
|
|
@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()})
|
|
|
|
|
|
async def _live_session_ids(client: ZertoClient) -> set[str]:
|
|
try:
|
|
rows = session_rows(await client.list_flrs())
|
|
except ZertoError:
|
|
return set()
|
|
return {session_id_of(r) for r in rows if is_live_session(r) and session_id_of(r)}
|
|
|
|
|
|
async def _teardown_flr(
|
|
client: ZertoClient,
|
|
session_id: str | None,
|
|
before: set[str],
|
|
) -> dict[str, Any]:
|
|
"""Unmount the FLR session and say what happened.
|
|
|
|
A swallowed unmount failure is how a mount silently wedges the recovery
|
|
site: FLR cannot run during clone/test/live/EJC, so a stuck session blocks
|
|
the next recovery. Report it instead.
|
|
|
|
If session_id is None the start may still have succeeded on the ZVM while
|
|
the response failed to parse, so reap anything new that appeared.
|
|
"""
|
|
out: dict[str, Any] = {"attempted": False, "ok": True, "ended": [], "failed": []}
|
|
targets = [session_id] if session_id else []
|
|
if not targets:
|
|
orphans = sorted(await _live_session_ids(client) - before)
|
|
targets = orphans
|
|
out["orphans_reaped"] = orphans
|
|
for sid in targets:
|
|
out["attempted"] = True
|
|
try:
|
|
await client.end_flr(sid)
|
|
out["ended"].append(sid)
|
|
except ZertoError as exc:
|
|
out["ok"] = False
|
|
out["failed"].append({"session_id": sid, "message": str(exc)})
|
|
if not out["ok"]:
|
|
out["message"] = (
|
|
"FLR session may still be mounted on the recovery site. "
|
|
"List it with zerto_list_flr_sessions and end it with "
|
|
"zerto_end_flr_session; a stuck mount blocks the next FLR."
|
|
)
|
|
return out
|
|
|
|
|
|
@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: str | None = None
|
|
before = await _live_session_ids(client)
|
|
payload: dict[str, Any]
|
|
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)
|
|
flr_path = await resolve_flr_path(client, session_id, guest_path)
|
|
token_payload = await client.download_flr(session_id, [flr_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)
|
|
payload = {
|
|
"ok": True,
|
|
"path": str(out_path.resolve()),
|
|
"bytes": len(blob),
|
|
"session_id": session_id,
|
|
"flr_path": flr_path,
|
|
"message": f"Wrote {len(blob)} bytes to {out_path}",
|
|
}
|
|
except ZertoError as exc:
|
|
payload = {"ok": False, "session_id": session_id, "message": str(exc)}
|
|
finally:
|
|
unmount = await _teardown_flr(client, session_id, before)
|
|
payload["unmount"] = unmount
|
|
if not unmount["ok"]:
|
|
# The bytes are on disk, so ok stays true, but the caller must be told
|
|
# the mount is still up rather than finding out at the next recovery.
|
|
payload["warning"] = unmount["message"]
|
|
return _dump(payload)
|
|
|
|
|
|
@mcp.tool(
|
|
name="zerto_list_flr_sessions",
|
|
annotations={
|
|
"title": "List FLR sessions on the ZVM",
|
|
"readOnlyHint": True,
|
|
"destructiveHint": False,
|
|
"idempotentHint": True,
|
|
"openWorldHint": True,
|
|
},
|
|
)
|
|
async def zerto_list_flr_sessions(live_only: bool = True) -> str:
|
|
"""Every FLR session the ZVM knows about, so orphaned mounts are visible.
|
|
|
|
zerto_recover_file tears its own session down, but that cleanup only runs if
|
|
this process survives the call. A crash, disconnect or timeout mid-recovery
|
|
leaves the mount up with nothing tracking it. FLR cannot run during clone,
|
|
test, live failover or EJC, so a stuck mount blocks the next recovery.
|
|
|
|
live_only keeps sessions still holding a mount. Pass false to see ended and
|
|
failed sessions too, which the ZVM keeps as history.
|
|
"""
|
|
try:
|
|
rows = session_rows(await get_client().list_flrs())
|
|
except ZertoError as exc:
|
|
return _dump({"ok": False, "message": str(exc)})
|
|
kept = [r for r in rows if is_live_session(r)] if live_only else rows
|
|
return _dump(
|
|
{
|
|
"ok": True,
|
|
"live_only": live_only,
|
|
"count": len(kept),
|
|
"total_known": len(rows),
|
|
"sessions": [session_summary(r) for r in kept],
|
|
}
|
|
)
|
|
|
|
|
|
@mcp.tool(
|
|
name="zerto_end_flr_session",
|
|
annotations={
|
|
"title": "End an FLR session (unmount)",
|
|
"readOnlyHint": False,
|
|
"destructiveHint": True,
|
|
"idempotentHint": True,
|
|
"openWorldHint": True,
|
|
},
|
|
)
|
|
async def zerto_end_flr_session(session_id: str, confirmed: bool = False) -> str:
|
|
"""Unmount an FLR session. Use to reap an orphan left by a crashed recovery.
|
|
|
|
Requires confirmed=true: ending a session that another operator is actively
|
|
pulling files from will interrupt them. Find the id with
|
|
zerto_list_flr_sessions.
|
|
"""
|
|
if not confirmed:
|
|
return _dump(
|
|
{
|
|
"ok": False,
|
|
"needs_confirm": True,
|
|
"message": (
|
|
"Set confirmed=true after a human yes. Ending a session that "
|
|
"someone is actively recovering from will interrupt them."
|
|
),
|
|
}
|
|
)
|
|
try:
|
|
await get_client().end_flr(session_id)
|
|
except ZertoError as exc:
|
|
return _dump({"ok": False, "session_id": session_id, "message": str(exc)})
|
|
return _dump(
|
|
{
|
|
"ok": True,
|
|
"session_id": session_id,
|
|
"message": f"Ended FLR session {session_id}.",
|
|
}
|
|
)
|
|
|
|
|
|
@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()
|