Compare commits

...
3 Commits
Author SHA1 Message Date
justinandClaude Opus 5 f930b84615 feat(flr): make FLR session lifecycle visible and reapable
zerto_recover_file already tore its session down in a finally block, but
three gaps meant a mount could stay up on the recovery site with nothing
tracking it. FLR cannot run during clone, test, live failover or EJC, so
a stuck session blocks the next recovery.

1. An unmount failure was swallowed (`except ZertoError: pass`). The
   caller got ok=true and never learned the mount was still up. The
   teardown result is now reported in the response as `unmount`, with a
   `warning` when it fails. ok stays true when the bytes did land -- the
   recovery genuinely succeeded -- but the caller is told.

2. If start_flr succeeded on the ZVM while its response failed to parse,
   session_id stayed None and the finally block did nothing, leaking a
   session the process never knew the id of. Teardown now snapshots live
   session ids before starting and reaps anything new that appeared,
   leaving other operators' sessions alone.

3. Nothing could see or clear an orphan left by a crashed process, since
   the finally block only runs if the process survives. Two new tools:

   - zerto_list_flr_sessions: every session the ZVM knows about.
     live_only (default true) keeps the ones still holding a mount;
     ended and failed sessions linger as history and hold nothing.
   - zerto_end_flr_session: unmount one. Gated on confirmed=true,
     matching the other destructive tools, because ending a session
     someone else is pulling files from will interrupt them.

Verified against ZVM 10.x: listing reports 0 live / 1 known after a clean
run, the confirm gate refuses without a human yes, a real recovery from
checkpoint 1368 returned 158 bytes and reported
unmount.ok=true with the session id it ended, and 0 live sessions
remained afterwards.

pytest 27 passed (5 new, including fakes covering the swallowed-failure
and orphan-reap paths).

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_016yVfC5nvZowoLFnEGWhLGn
2026-09-21 13:40:53 -04:00
justinandClaude Opus 5 fe7b220ab9 feat(checkpoints): record agent and intent in the checkpoint name
Zerto's tagged checkpoint insert takes exactly one field. The 10.x
swagger model VpgInsertTagCheckpointDataApi has a single property,
checkpointName, and the 9.0 API reference lists CheckpointName as the
only request value. There is no description field, so who the agent is
and what it is about to do have to live inside the name.

Old name:
  ai:claude:chg-412:20260921T170829Z

New name:
  ai:claude | edit /home/justin/app-config.yaml | vm=jp-ubuntu |
  change=chg-412 | 20260921T170829Z

zerto_create_tagged_checkpoint and zerto_guard_before_mutate take a new
action argument: free text saying what the agent is about to do. The VM
name is filled in from the find result. An operator reading the journal
in the Zerto UI can now see which agent inserted a checkpoint and why,
without the agent transcript.

Field text is sanitised so the name stays one readable line: control
characters and runs of whitespace collapse to single spaces, ';' becomes
',' because Zerto appends "; Used for File Level Restore" to its own
tags, and '|' becomes '/' because ' | ' is our field separator. Capped
at TAG_MAX_LEN (250).

Measured against ZVM 10.x while picking the format:

- names of at least 400 chars are accepted, and spaces, slashes,
  parentheses, '=' and '|' all survive the round trip
- tagged checkpoint inserts fired back to back at one VPG are silently
  dropped. The POST returns 200 and queues a task, but only the first
  checkpoint appears. tag_vpgs already inserts then waits per VPG, so
  it is correct; added a comment so nobody turns that loop into an
  asyncio.gather().

Verified end to end: guard inserted cp 1197 on VPG jp-ubuntu, the name
read back byte-identical from the journal, and FLR from that checkpoint
returned the 158 byte pre-mutation file.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_016yVfC5nvZowoLFnEGWhLGn
2026-09-21 13:09:08 -04:00
justinandClaude Opus 5 94025b8d69 fix(flr): resolve guest paths into the FLR partition namespace
zerto_recover_file passed the guest absolute path straight to
POST /v1/flrs/{session}/download, which the ZVM rejects:

  HTTP 400 {"Message":"Invalid path: check location exists or
  correct path syntax."}

FLR browse/download is rooted at partitions, not the guest's /.
/home/justin/app-config.yaml is Volume2-Ext4/home/justin/app-config.yaml.
server.py already called browse_flr() but discarded the result, so
nothing ever resolved the path.

Add resolve_flr_path() and browsable_partitions() to recover.py:

- browse path "" returns {MainPathItem, PathItems}, not a bare list
- skip partitions with IsBrowsable false. A Linux guest reports
  Volume1-Unknown as "Cannot restore. Partition type Unknown is not
  supported." Do not assume the first partition is the right one.
- browse returns child paths percent-encoded
  (Volume2-Ext4%2fhome%2fjustin%2fapp-config.yaml); download wants
  them decoded with plain slashes
- raise a ZertoError naming what was searched when the file is
  absent, since "not replicated into that checkpoint yet" is the
  likely cause and is actionable

zerto_recover_file now returns the resolved flr_path.

Verified end to end against ZVM 10.x: guard tagged cp 1075 on VPG
jp-ubuntu, FLR mounted in ~3s, 158 bytes recovered from the
pre-mutation checkpoint with matching content.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_016yVfC5nvZowoLFnEGWhLGn
2026-09-21 13:04:08 -04:00
10 changed files with 530 additions and 45 deletions
+1 -1
View File
@@ -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**:
+2 -2
View File
@@ -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.
+5
View File
@@ -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.
+20 -2
View File
@@ -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.
+49 -5
View File
@@ -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] = []
+4
View File
@@ -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}")
+103
View File
@@ -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
+178 -26
View File
@@ -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(
try: {
await client.end_flr(session_id) "ok": True,
except ZertoError: "live_only": live_only,
pass "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( @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."
), ),
} }
) )
+37 -9
View File
@@ -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():
+131
View File
@@ -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 == []