Compare commits
3
Commits
main
...
f930b84615
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f930b84615 | ||
|
|
fe7b220ab9 | ||
|
|
94025b8d69 |
+1
-1
@@ -5,7 +5,7 @@ PoC MCP that teaches an agent to discover Zerto protection, pin a tagged checkpo
|
|||||||
## Language
|
## Language
|
||||||
|
|
||||||
**Tagged checkpoint**:
|
**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.
|
A named bookmark in a VPG journal, inserted by `POST /v1/vpgs/{id}/checkpoints` (`startVpgTaggedCheckpointInsert`). Crash-consistent write-order only; not application-quiesced unless someone scripted that separately. `CheckpointName` is the only field the API accepts, so agent and intent go in the name: `ai:<agent> | <action> | vm=<vm> | change=<id> | <utc>`. Inserts are async tasks and are silently dropped if fired back to back at one VPG; insert, then wait until listed.
|
||||||
_Avoid_: user checkpoint, snapshot, backup, restore point (unqualified)
|
_Avoid_: user checkpoint, snapshot, backup, restore point (unqualified)
|
||||||
|
|
||||||
**VPG**:
|
**VPG**:
|
||||||
|
|||||||
@@ -7,8 +7,8 @@ If the loop works, these tools are the delta to add to official ZVM MCP (`ZVM.MC
|
|||||||
## What it does
|
## 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.
|
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.
|
2. `zerto_create_tagged_checkpoint` / `zerto_guard_before_mutate` — same tag on every protecting VPG, wait until listed. The name records which agent and what it is doing: `ai:<agent> | <action> | vm=<vm> | change=<id> | <utc>`.
|
||||||
3. `zerto_recover_file` — FLR after a human sets `confirmed=true`.
|
3. `zerto_recover_file` — FLR after a human sets `confirmed=true`. Reports its own unmount; `zerto_list_flr_sessions` / `zerto_end_flr_session` find and reap a mount orphaned by a crashed recovery.
|
||||||
4. Mutating catalog — opt-in list of MCP tools that must be guarded. Unlisted tools pass through. Users add entries.
|
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.
|
Official ZVM MCP already has inventory and failover test. It does not insert tagged checkpoints or run FLR.
|
||||||
|
|||||||
@@ -24,6 +24,11 @@ Do not use FLR when:
|
|||||||
- OS-level dedup volumes
|
- OS-level dedup volumes
|
||||||
- 10.9 FLR Operator role (broken; Administrator is the documented workaround)
|
- 10.9 FLR Operator role (broken; Administrator is the documented workaround)
|
||||||
|
|
||||||
|
FLR sessions must be unmounted when done. `zerto_recover_file` does that
|
||||||
|
itself and reports it in `unmount`, but that cleanup only runs if the MCP
|
||||||
|
process survives the call. After a crash, list orphans with
|
||||||
|
`zerto_list_flr_sessions` and end them with `zerto_end_flr_session`.
|
||||||
|
|
||||||
A tagged checkpoint must already exist. Initial sync has an empty journal
|
A tagged checkpoint must already exist. Initial sync has an empty journal
|
||||||
(`GET .../checkpoints` returns `[]`). Guard refuses until status is MeetingSLA
|
(`GET .../checkpoints` returns `[]`). Guard refuses until status is MeetingSLA
|
||||||
(or NotMeetingSLA) and substatus is not a sync.
|
(or NotMeetingSLA) and substatus is not a sync.
|
||||||
|
|||||||
@@ -14,7 +14,7 @@ You talk to **one** MCP: `zerto_rewind_mcp`. Do not also require official ZVM MC
|
|||||||
Before **every** guest-mutating tool call:
|
Before **every** guest-mutating tool call:
|
||||||
|
|
||||||
1. Take the hostname / VM name / Zerto `vmIdentifier` from the tool args.
|
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`).
|
2. Call `zerto_guard_before_mutate` with `change_id` and `action` (or `zerto_find_protection` then `zerto_create_tagged_checkpoint`).
|
||||||
3. If `ok` is not true: **stop**. Do not mutate.
|
3. If `ok` is not true: **stop**. Do not mutate.
|
||||||
4. Then run the mutating call.
|
4. Then run the mutating call.
|
||||||
|
|
||||||
@@ -41,6 +41,12 @@ Human must confirm. Pass `confirmed=true` only after they say yes.
|
|||||||
- Inspect a whole VM: `zerto_offsite_clone` or `zerto_start_failover_test`.
|
- Inspect a whole VM: `zerto_offsite_clone` or `zerto_start_failover_test`.
|
||||||
- Never Failover Live. Never Move. Those are DR, not rewind.
|
- Never Failover Live. Never Move. Those are DR, not rewind.
|
||||||
|
|
||||||
|
`zerto_recover_file` unmounts its own FLR session and reports the result in
|
||||||
|
`unmount`. If `unmount.ok` is false, or a previous recovery died mid-flight,
|
||||||
|
the mount is still up: FLR cannot run during clone, test, live failover or EJC,
|
||||||
|
so a stuck session blocks the next recovery. Find it with
|
||||||
|
`zerto_list_flr_sessions` and clear it with `zerto_end_flr_session`.
|
||||||
|
|
||||||
## Facts that bite
|
## Facts that bite
|
||||||
|
|
||||||
- A tagged checkpoint is crash-consistent, not app-quiesced.
|
- A tagged checkpoint is crash-consistent, not app-quiesced.
|
||||||
@@ -51,4 +57,16 @@ Human must confirm. Pass `confirmed=true` only after they say yes.
|
|||||||
|
|
||||||
## Tag
|
## Tag
|
||||||
|
|
||||||
Default: `ai:<agent>:<change-id>:<utc>`. Same string on every VPG for that call.
|
The checkpoint name is the only field the Zerto API takes, so it carries the
|
||||||
|
whole story:
|
||||||
|
|
||||||
|
```
|
||||||
|
ai:<agent> | <action> | vm=<vm> | change=<change-id> | <utc>
|
||||||
|
ai:claude | edit /etc/nginx/nginx.conf | vm=web01 | change=chg-412 | 20260921T150405Z
|
||||||
|
```
|
||||||
|
|
||||||
|
Always pass `action`: a plain description of the change you are about to make.
|
||||||
|
An operator scrolling the journal in the Zerto UI should be able to tell which
|
||||||
|
agent inserted the checkpoint and why, without reading your transcript.
|
||||||
|
|
||||||
|
Same string on every VPG for that call.
|
||||||
|
|||||||
@@ -11,14 +11,54 @@ from zerto_rewind_mcp.client import ZertoClient, ZertoError
|
|||||||
from zerto_rewind_mcp.protection import FindResult
|
from zerto_rewind_mcp.protection import FindResult
|
||||||
from zerto_rewind_mcp.util import pick
|
from zerto_rewind_mcp.util import pick
|
||||||
|
|
||||||
_SAFE = re.compile(r"[^A-Za-z0-9._:-]+")
|
_CTRL = re.compile(r"[\x00-\x1f\x7f]+")
|
||||||
|
_WS = re.compile(r"\s+")
|
||||||
|
|
||||||
|
# CheckpointName is the only field POST /v1/vpgs/{id}/checkpoints accepts
|
||||||
|
# (VpgInsertTagCheckpointDataApi has exactly one property). Who the agent is
|
||||||
|
# and what it is about to do therefore have to live inside the name itself.
|
||||||
|
# Measured on 10.x: names of at least 400 chars are accepted, and spaces,
|
||||||
|
# slashes, parentheses, '=' and '|' all survive. We still cap at TAG_MAX_LEN
|
||||||
|
# so the name stays readable in the Zerto UI checkpoint list.
|
||||||
|
TAG_MAX_LEN = 250
|
||||||
|
|
||||||
|
|
||||||
def make_tag(agent: str, change_id: str, when: datetime | None = None) -> str:
|
def _clean(value: Any, limit: int) -> str:
|
||||||
|
"""One field of a checkpoint name: single-line, no separator collisions."""
|
||||||
|
text = _CTRL.sub(" ", str(value or ""))
|
||||||
|
# Zerto appends "; Used for File Level Restore" to its own tags, and we use
|
||||||
|
# " | " as our field separator. Keep both out of user-supplied text.
|
||||||
|
text = text.replace(";", ",").replace("|", "/")
|
||||||
|
return _WS.sub(" ", text).strip()[:limit]
|
||||||
|
|
||||||
|
|
||||||
|
def make_tag(
|
||||||
|
agent: str,
|
||||||
|
change_id: str,
|
||||||
|
when: datetime | None = None,
|
||||||
|
*,
|
||||||
|
action: str | None = None,
|
||||||
|
vm_name: str | None = None,
|
||||||
|
) -> str:
|
||||||
|
"""Build the checkpoint name.
|
||||||
|
|
||||||
|
Says which agent, what it is about to do, to which VM, under what change id:
|
||||||
|
|
||||||
|
ai:claude | edit /etc/app.conf | vm=jp-ubuntu | change=chg-99 | 20260921T150405Z
|
||||||
|
|
||||||
|
action is free text from the caller describing the pending mutation.
|
||||||
|
"""
|
||||||
stamp = (when or datetime.now(UTC)).strftime("%Y%m%dT%H%M%SZ")
|
stamp = (when or datetime.now(UTC)).strftime("%Y%m%dT%H%M%SZ")
|
||||||
agent_s = _SAFE.sub("-", agent.strip())[:40] or "agent"
|
parts = [f"ai:{_clean(agent, 40) or 'agent'}"]
|
||||||
change_s = _SAFE.sub("-", change_id.strip())[:80] or "change"
|
cleaned_action = _clean(action, 120)
|
||||||
return f"ai:{agent_s}:{change_s}:{stamp}"
|
if cleaned_action:
|
||||||
|
parts.append(cleaned_action)
|
||||||
|
cleaned_vm = _clean(vm_name, 60)
|
||||||
|
if cleaned_vm:
|
||||||
|
parts.append(f"vm={cleaned_vm}")
|
||||||
|
parts.append(f"change={_clean(change_id, 80) or 'change'}")
|
||||||
|
parts.append(stamp)
|
||||||
|
return " | ".join(parts)[:TAG_MAX_LEN]
|
||||||
|
|
||||||
|
|
||||||
def checkpoint_tag(row: dict[str, Any]) -> str:
|
def checkpoint_tag(row: dict[str, Any]) -> str:
|
||||||
@@ -68,6 +108,10 @@ async def tag_vpgs(
|
|||||||
"message": result.message,
|
"message": result.message,
|
||||||
"find": result.as_dict(),
|
"find": result.as_dict(),
|
||||||
}
|
}
|
||||||
|
# Insert then wait, one VPG at a time. Measured on 10.x: tagged checkpoint
|
||||||
|
# inserts fired back to back at the same VPG are silently dropped -- the POST
|
||||||
|
# returns 200 and queues a task, but only the first checkpoint ever appears.
|
||||||
|
# Do not turn this loop into an asyncio.gather().
|
||||||
tagged: list[dict[str, Any]] = []
|
tagged: list[dict[str, Any]] = []
|
||||||
skipped = [v.as_dict() for v in result.vm.vpgs if not v.can_tag]
|
skipped = [v.as_dict() for v in result.vm.vpgs if not v.can_tag]
|
||||||
errors: list[str] = []
|
errors: list[str] = []
|
||||||
|
|||||||
@@ -231,6 +231,10 @@ class ZertoClient:
|
|||||||
status_code=last.status_code if last else None,
|
status_code=last.status_code if last else None,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
async def list_flrs(self) -> Any:
|
||||||
|
"""GET /v1/flrs. Every FLR session the ZVM currently knows about."""
|
||||||
|
return await self.json("GET", "/v1/flrs")
|
||||||
|
|
||||||
async def end_flr(self, session_id: str) -> None:
|
async def end_flr(self, session_id: str) -> None:
|
||||||
try:
|
try:
|
||||||
await self.json("DELETE", f"/v1/flrs/{session_id}")
|
await self.json("DELETE", f"/v1/flrs/{session_id}")
|
||||||
|
|||||||
@@ -3,6 +3,7 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import urllib.parse
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from zerto_rewind_mcp.client import ZertoClient, ZertoError
|
from zerto_rewind_mcp.client import ZertoClient, ZertoError
|
||||||
@@ -94,3 +95,105 @@ async def wait_flr_ready(
|
|||||||
f"FLR session {session_id} not ready within {timeout_s:.0f}s (last={last!r}). "
|
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."
|
"Retry browse after a minute, or check EJC is not running and no clone/test is active."
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _decode_flr_path(value: Any) -> str:
|
||||||
|
"""Browse returns paths percent-encoded (%2f). Download wants them decoded."""
|
||||||
|
return urllib.parse.unquote(str(value or "")).replace("\\", "/")
|
||||||
|
|
||||||
|
|
||||||
|
def path_items(payload: Any) -> list[dict[str, Any]]:
|
||||||
|
if isinstance(payload, dict):
|
||||||
|
items = payload.get("PathItems") or payload.get("pathItems") or []
|
||||||
|
return [i for i in items if isinstance(i, dict)]
|
||||||
|
if isinstance(payload, list):
|
||||||
|
return [i for i in payload if isinstance(i, dict)]
|
||||||
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
async def browsable_partitions(client: ZertoClient, session_id: str) -> list[str]:
|
||||||
|
"""FLR is rooted at partitions (Volume2-Ext4), not the guest's /.
|
||||||
|
|
||||||
|
Volume1-Unknown and friends report IsBrowsable false and cannot be restored.
|
||||||
|
"""
|
||||||
|
rows = path_items(await client.browse_flr(session_id, path=""))
|
||||||
|
return [str(r.get("Path")) for r in rows if r.get("IsBrowsable")]
|
||||||
|
|
||||||
|
|
||||||
|
async def resolve_flr_path(client: ZertoClient, session_id: str, guest_path: str) -> str:
|
||||||
|
"""Map a guest absolute path to the FLR namespace path the download API accepts.
|
||||||
|
|
||||||
|
/home/justin/app-config.yaml -> Volume2-Ext4/home/justin/app-config.yaml
|
||||||
|
"""
|
||||||
|
rel = guest_path.replace("\\", "/").strip("/")
|
||||||
|
if not rel:
|
||||||
|
raise ZertoError("Empty guest_path")
|
||||||
|
parts = rel.split("/")
|
||||||
|
name, parent = parts[-1], "/".join(parts[:-1])
|
||||||
|
|
||||||
|
partitions = await browsable_partitions(client, session_id)
|
||||||
|
if not partitions:
|
||||||
|
raise ZertoError(
|
||||||
|
"FLR mounted but no browsable partition. Unsupported partition type "
|
||||||
|
"(LVM/unknown) cannot be restored by FLR."
|
||||||
|
)
|
||||||
|
tried = []
|
||||||
|
for vol in partitions:
|
||||||
|
probe = f"{vol}/{parent}" if parent else vol
|
||||||
|
tried.append(probe)
|
||||||
|
try:
|
||||||
|
rows = path_items(await client.browse_flr(session_id, path=probe))
|
||||||
|
except ZertoError:
|
||||||
|
continue
|
||||||
|
for row in rows:
|
||||||
|
decoded = _decode_flr_path(row.get("Path"))
|
||||||
|
if decoded.rsplit("/", 1)[-1] == name:
|
||||||
|
return decoded
|
||||||
|
raise ZertoError(
|
||||||
|
f"{guest_path!r} not found in the FLR mount. Looked under {tried}. "
|
||||||
|
"The file may not have replicated into that checkpoint yet."
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def session_rows(payload: Any) -> list[dict[str, Any]]:
|
||||||
|
"""Normalise GET /v1/flrs, which returns a list or a single object."""
|
||||||
|
if isinstance(payload, list):
|
||||||
|
return [r for r in payload if isinstance(r, dict)]
|
||||||
|
if isinstance(payload, dict):
|
||||||
|
return [payload]
|
||||||
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
def session_id_of(row: dict[str, Any]) -> str:
|
||||||
|
value = pick(
|
||||||
|
row,
|
||||||
|
"FlrSessionIdentifier",
|
||||||
|
"flrSessionIdentifier",
|
||||||
|
"SessionId",
|
||||||
|
"sessionId",
|
||||||
|
"Identifier",
|
||||||
|
"identifier",
|
||||||
|
)
|
||||||
|
return str(value) if value is not None else ""
|
||||||
|
|
||||||
|
|
||||||
|
def session_summary(row: dict[str, Any]) -> dict[str, Any]:
|
||||||
|
return {
|
||||||
|
"session_id": session_id_of(row),
|
||||||
|
"state": flr_status(row),
|
||||||
|
"vpg_name": pick(row, "VpgName", "vpgName"),
|
||||||
|
"vm_name": pick(row, "VmName", "vmName"),
|
||||||
|
"checkpoint_id": pick(row, "CheckpointIdentifier", "checkpointIdentifier"),
|
||||||
|
"mounted_at": pick(row, "MountedTime", "mountedTime", "StartTime", "startTime"),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def is_live_session(row: dict[str, Any]) -> bool:
|
||||||
|
"""A session still holding a mount on the recovery site.
|
||||||
|
|
||||||
|
Unmounted/ended sessions linger in GET /v1/flrs as history; they hold nothing.
|
||||||
|
"""
|
||||||
|
state = flr_status(row).lower()
|
||||||
|
if not state:
|
||||||
|
return False
|
||||||
|
return "unmount" not in state and "fail" not in state and "end" not in state
|
||||||
|
|||||||
+171
-19
@@ -13,7 +13,16 @@ from zerto_rewind_mcp.checkpoints import checkpoint_id, checkpoint_tag, make_tag
|
|||||||
from zerto_rewind_mcp.client import ZertoClient, ZertoError
|
from zerto_rewind_mcp.client import ZertoClient, ZertoError
|
||||||
from zerto_rewind_mcp.config import load_catalog, load_config
|
from zerto_rewind_mcp.config import load_catalog, load_config
|
||||||
from zerto_rewind_mcp.protection import find_from_rows
|
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.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
|
from zerto_rewind_mcp.util import pick
|
||||||
|
|
||||||
mcp = FastMCP("zerto_rewind_mcp")
|
mcp = FastMCP("zerto_rewind_mcp")
|
||||||
@@ -110,12 +119,20 @@ async def zerto_create_tagged_checkpoint(
|
|||||||
query: str,
|
query: str,
|
||||||
change_id: str,
|
change_id: str,
|
||||||
agent: str = "agent",
|
agent: str = "agent",
|
||||||
|
action: str = "",
|
||||||
tag: str | None = None,
|
tag: str | None = None,
|
||||||
) -> str:
|
) -> str:
|
||||||
"""Insert the same tagged checkpoint on every protecting VPG, then wait.
|
"""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.
|
Call this before every cataloged guest-mutating tool. If it fails, refuse the change.
|
||||||
Tag format defaults to ai:<agent>:<change_id>:<utc>.
|
|
||||||
|
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;
|
Docs: tagged checkpoints are not supported when the protected site is Azure or AWS;
|
||||||
run this against the vSphere protected ZVM.
|
run this against the vSphere protected ZVM.
|
||||||
"""
|
"""
|
||||||
@@ -125,7 +142,12 @@ async def zerto_create_tagged_checkpoint(
|
|||||||
payload = result.as_dict()
|
payload = result.as_dict()
|
||||||
payload["ok"] = False
|
payload["ok"] = False
|
||||||
return _dump(payload)
|
return _dump(payload)
|
||||||
use_tag = tag or make_tag(agent, change_id)
|
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)
|
out = await tag_vpgs(get_client(), result, use_tag)
|
||||||
return _dump(out)
|
return _dump(out)
|
||||||
except ZertoError as exc:
|
except ZertoError as exc:
|
||||||
@@ -182,13 +204,19 @@ async def zerto_guard_before_mutate(
|
|||||||
query: str,
|
query: str,
|
||||||
change_id: str,
|
change_id: str,
|
||||||
agent: str = "agent",
|
agent: str = "agent",
|
||||||
|
action: str = "",
|
||||||
) -> str:
|
) -> str:
|
||||||
"""Find protection and tag. If ok is not true, do not mutate.
|
"""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:
|
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.
|
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)
|
return await zerto_create_tagged_checkpoint(
|
||||||
|
query=query, change_id=change_id, agent=agent, action=action
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@mcp.tool(
|
@mcp.tool(
|
||||||
@@ -232,6 +260,51 @@ async def zerto_add_mutating_tool(
|
|||||||
return _dump({"ok": True, "entry": entry.as_dict()})
|
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(
|
@mcp.tool(
|
||||||
name="zerto_recover_file",
|
name="zerto_recover_file",
|
||||||
annotations={
|
annotations={
|
||||||
@@ -261,15 +334,16 @@ async def zerto_recover_file(
|
|||||||
"ok": False,
|
"ok": False,
|
||||||
"needs_confirm": True,
|
"needs_confirm": True,
|
||||||
"message": (
|
"message": (
|
||||||
"Set confirmed=true after a human yes. "
|
"Set confirmed=true after a human yes. FLR mounts a disk on the recovery site."
|
||||||
"FLR mounts a disk on the recovery site."
|
|
||||||
),
|
),
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
client = get_client()
|
client = get_client()
|
||||||
dest = Path(dest_dir or _settings.get("recovery_dir") or "./recovered")
|
dest = Path(dest_dir or _settings.get("recovery_dir") or "./recovered")
|
||||||
dest.mkdir(parents=True, exist_ok=True)
|
dest.mkdir(parents=True, exist_ok=True)
|
||||||
session_id = None
|
session_id: str | None = None
|
||||||
|
before = await _live_session_ids(client)
|
||||||
|
payload: dict[str, Any]
|
||||||
try:
|
try:
|
||||||
started = await client.start_flr(
|
started = await client.start_flr(
|
||||||
vpg_identifier,
|
vpg_identifier,
|
||||||
@@ -279,30 +353,109 @@ async def zerto_recover_file(
|
|||||||
)
|
)
|
||||||
session_id = session_id_from(started)
|
session_id = session_id_from(started)
|
||||||
await wait_flr_ready(client, session_id)
|
await wait_flr_ready(client, session_id)
|
||||||
await client.browse_flr(session_id, path="", recursive=False)
|
flr_path = await resolve_flr_path(client, session_id, guest_path)
|
||||||
token_payload = await client.download_flr(session_id, [guest_path])
|
token_payload = await client.download_flr(session_id, [flr_path])
|
||||||
token = download_token_from(token_payload)
|
token = download_token_from(token_payload)
|
||||||
blob = await client.fetch_download(token)
|
blob = await client.fetch_download(token)
|
||||||
name = Path(guest_path.replace("\\", "/")).name or "recovered.bin"
|
name = Path(guest_path.replace("\\", "/")).name or "recovered.bin"
|
||||||
out_path = dest / name
|
out_path = dest / name
|
||||||
out_path.write_bytes(blob)
|
out_path.write_bytes(blob)
|
||||||
return _dump(
|
payload = {
|
||||||
{
|
|
||||||
"ok": True,
|
"ok": True,
|
||||||
"path": str(out_path.resolve()),
|
"path": str(out_path.resolve()),
|
||||||
"bytes": len(blob),
|
"bytes": len(blob),
|
||||||
"session_id": session_id,
|
"session_id": session_id,
|
||||||
|
"flr_path": flr_path,
|
||||||
"message": f"Wrote {len(blob)} bytes to {out_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:
|
except ZertoError as exc:
|
||||||
return _dump({"ok": False, "message": str(exc)})
|
return _dump({"ok": False, "message": str(exc)})
|
||||||
finally:
|
kept = [r for r in rows if is_live_session(r)] if live_only else rows
|
||||||
if session_id:
|
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:
|
try:
|
||||||
await client.end_flr(session_id)
|
await get_client().end_flr(session_id)
|
||||||
except ZertoError:
|
except ZertoError as exc:
|
||||||
pass
|
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(
|
@mcp.tool(
|
||||||
@@ -394,8 +547,7 @@ async def zerto_offsite_clone(
|
|||||||
"ok": False,
|
"ok": False,
|
||||||
"needs_confirm": True,
|
"needs_confirm": True,
|
||||||
"message": (
|
"message": (
|
||||||
"Set confirmed=true after a human yes. "
|
"Set confirmed=true after a human yes. Clone consumes recovery datastore."
|
||||||
"Clone consumes recovery datastore."
|
|
||||||
),
|
),
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -1,18 +1,46 @@
|
|||||||
from datetime import UTC, datetime
|
from datetime import UTC, datetime
|
||||||
|
|
||||||
from zerto_rewind_mcp.checkpoints import checkpoint_id, checkpoint_tag, make_tag
|
from zerto_rewind_mcp.checkpoints import (
|
||||||
|
TAG_MAX_LEN,
|
||||||
|
checkpoint_id,
|
||||||
|
checkpoint_tag,
|
||||||
|
make_tag,
|
||||||
|
)
|
||||||
|
|
||||||
|
WHEN = datetime(2026, 9, 21, 15, 4, 5, tzinfo=UTC)
|
||||||
|
|
||||||
|
|
||||||
def test_tag_format():
|
def test_tag_describes_agent_and_action():
|
||||||
when = datetime(2026, 9, 21, 15, 4, 5, tzinfo=UTC)
|
tag = make_tag(
|
||||||
tag = make_tag("codex", "chg-99", when)
|
"claude",
|
||||||
assert tag == "ai:codex:chg-99:20260921T150405Z"
|
"chg-99",
|
||||||
|
WHEN,
|
||||||
|
action="edit /home/justin/app-config.yaml",
|
||||||
|
vm_name="jp-ubuntu",
|
||||||
|
)
|
||||||
|
assert tag == (
|
||||||
|
"ai:claude | edit /home/justin/app-config.yaml | vm=jp-ubuntu "
|
||||||
|
"| change=chg-99 | 20260921T150405Z"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def test_tag_strips_junk():
|
def test_tag_without_action_or_vm():
|
||||||
tag = make_tag("agent one", "path/with spaces", datetime(2026, 1, 1, tzinfo=UTC))
|
assert make_tag("codex", "chg-99", WHEN) == "ai:codex | change=chg-99 | 20260921T150405Z"
|
||||||
assert " " not in tag
|
|
||||||
assert tag.startswith("ai:agent-one:path-with-spaces:")
|
|
||||||
|
def test_tag_keeps_separators_unambiguous():
|
||||||
|
# ';' is what Zerto appends to its own tags; '|' is our field separator.
|
||||||
|
tag = make_tag("a;b", "c|d", WHEN, action="rm -rf /tmp;x", vm_name="v|m")
|
||||||
|
assert tag.count(" | ") == 4
|
||||||
|
assert ";" not in tag
|
||||||
|
|
||||||
|
|
||||||
|
def test_tag_collapses_whitespace_and_caps_length():
|
||||||
|
tag = make_tag("agent one", "chg 1", WHEN, action="do\n many things")
|
||||||
|
assert "\n" not in tag
|
||||||
|
assert "do many things" in tag
|
||||||
|
long_tag = make_tag("a" * 200, "b" * 200, WHEN, action="c" * 400)
|
||||||
|
assert len(long_tag) <= TAG_MAX_LEN
|
||||||
|
|
||||||
|
|
||||||
def test_checkpoint_row_keys():
|
def test_checkpoint_row_keys():
|
||||||
|
|||||||
@@ -28,3 +28,134 @@ def test_flr_list_status():
|
|||||||
]
|
]
|
||||||
row = flr_row(payload)
|
row = flr_row(payload)
|
||||||
assert flr_status(row).lower() == "mountcompletedsuccessfully"
|
assert flr_status(row).lower() == "mountcompletedsuccessfully"
|
||||||
|
|
||||||
|
|
||||||
|
def test_path_items_shapes():
|
||||||
|
from zerto_rewind_mcp.recover import path_items
|
||||||
|
|
||||||
|
assert path_items({"PathItems": [{"Path": "a"}]}) == [{"Path": "a"}]
|
||||||
|
assert path_items([{"Path": "b"}]) == [{"Path": "b"}]
|
||||||
|
assert path_items(None) == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_decode_flr_path():
|
||||||
|
from zerto_rewind_mcp.recover import _decode_flr_path
|
||||||
|
|
||||||
|
assert _decode_flr_path("Volume2-Ext4%2fhome%2fjustin%2fapp-config.yaml") == (
|
||||||
|
"Volume2-Ext4/home/justin/app-config.yaml"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_resolve_flr_path_picks_browsable_partition():
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from zerto_rewind_mcp.recover import resolve_flr_path
|
||||||
|
|
||||||
|
root = {
|
||||||
|
"PathItems": [
|
||||||
|
{"Path": "Volume1-Unknown", "IsBrowsable": False},
|
||||||
|
{"Path": "Volume2-Ext4", "IsBrowsable": True},
|
||||||
|
]
|
||||||
|
}
|
||||||
|
listing = {
|
||||||
|
"PathItems": [
|
||||||
|
{"Path": "Volume2-Ext4%2fhome%2fjustin%2f.bashrc", "Type": "File"},
|
||||||
|
{"Path": "Volume2-Ext4%2fhome%2fjustin%2fapp-config.yaml", "Type": "File"},
|
||||||
|
]
|
||||||
|
}
|
||||||
|
|
||||||
|
class FakeClient:
|
||||||
|
def __init__(self):
|
||||||
|
self.seen = []
|
||||||
|
|
||||||
|
async def browse_flr(self, session_id, path="", recursive=False):
|
||||||
|
self.seen.append(path)
|
||||||
|
return root if path == "" else listing
|
||||||
|
|
||||||
|
client = FakeClient()
|
||||||
|
got = asyncio.run(resolve_flr_path(client, "sess", "/home/justin/app-config.yaml"))
|
||||||
|
assert got == "Volume2-Ext4/home/justin/app-config.yaml"
|
||||||
|
# must not try the unrestorable partition
|
||||||
|
assert "Volume1-Unknown/home/justin" not in client.seen
|
||||||
|
|
||||||
|
|
||||||
|
def test_session_rows_and_id_shapes():
|
||||||
|
from zerto_rewind_mcp.recover import session_id_of, session_rows
|
||||||
|
|
||||||
|
assert session_rows([{"a": 1}]) == [{"a": 1}]
|
||||||
|
assert session_rows({"a": 1}) == [{"a": 1}]
|
||||||
|
assert session_rows(None) == []
|
||||||
|
assert session_id_of({"FlrSessionIdentifier": "s1"}) == "s1"
|
||||||
|
assert session_id_of({"sessionId": "s2"}) == "s2"
|
||||||
|
assert session_id_of({}) == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_is_live_session():
|
||||||
|
from zerto_rewind_mcp.recover import is_live_session
|
||||||
|
|
||||||
|
assert is_live_session({"FlrSessionStatus": "MountCompletedSuccessfully"})
|
||||||
|
assert is_live_session({"FlrSessionStatus": "MountInProgress"})
|
||||||
|
# unmounted/ended/failed sessions linger as history and hold nothing
|
||||||
|
assert not is_live_session({"FlrSessionStatus": "UnmountCompletedSuccessfully"})
|
||||||
|
assert not is_live_session({"FlrSessionStatus": "MountFailed"})
|
||||||
|
assert not is_live_session({})
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeClient:
|
||||||
|
"""Stands in for ZertoClient in teardown tests."""
|
||||||
|
|
||||||
|
def __init__(self, sessions=None, fail_end=False):
|
||||||
|
self._sessions = sessions or []
|
||||||
|
self.fail_end = fail_end
|
||||||
|
self.ended: list[str] = []
|
||||||
|
|
||||||
|
async def list_flrs(self):
|
||||||
|
return self._sessions
|
||||||
|
|
||||||
|
async def end_flr(self, session_id):
|
||||||
|
from zerto_rewind_mcp.client import ZertoError
|
||||||
|
|
||||||
|
if self.fail_end:
|
||||||
|
raise ZertoError("unmount refused")
|
||||||
|
self.ended.append(session_id)
|
||||||
|
|
||||||
|
|
||||||
|
def test_teardown_reports_failure_instead_of_swallowing():
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from zerto_rewind_mcp.server import _teardown_flr
|
||||||
|
|
||||||
|
client = _FakeClient(fail_end=True)
|
||||||
|
out = asyncio.run(_teardown_flr(client, "sess-1", set()))
|
||||||
|
assert out["ok"] is False
|
||||||
|
assert out["failed"][0]["session_id"] == "sess-1"
|
||||||
|
assert "still be mounted" in out["message"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_teardown_reaps_orphan_when_session_id_never_parsed():
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from zerto_rewind_mcp.server import _teardown_flr
|
||||||
|
|
||||||
|
# start_flr succeeded on the ZVM but the response did not parse, so the
|
||||||
|
# caller never learned the id. The new live session must still be reaped.
|
||||||
|
live = [{"FlrSessionIdentifier": "new-1", "FlrSessionStatus": "MountCompletedSuccessfully"}]
|
||||||
|
client = _FakeClient(sessions=live)
|
||||||
|
out = asyncio.run(_teardown_flr(client, None, set()))
|
||||||
|
assert out["ended"] == ["new-1"]
|
||||||
|
assert out["orphans_reaped"] == ["new-1"]
|
||||||
|
assert client.ended == ["new-1"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_teardown_leaves_pre_existing_sessions_alone():
|
||||||
|
import asyncio
|
||||||
|
|
||||||
|
from zerto_rewind_mcp.server import _teardown_flr
|
||||||
|
|
||||||
|
live = [
|
||||||
|
{"FlrSessionIdentifier": "someone-else", "FlrSessionStatus": "MountCompletedSuccessfully"}
|
||||||
|
]
|
||||||
|
client = _FakeClient(sessions=live)
|
||||||
|
out = asyncio.run(_teardown_flr(client, None, {"someone-else"}))
|
||||||
|
assert out["ended"] == []
|
||||||
|
assert client.ended == []
|
||||||
|
|||||||
Reference in New Issue
Block a user