fix(flr): partition-rooted FLR paths + agent/intent in checkpoint names (#2)

This commit was merged in pull request #2.
This commit is contained in:
2026-09-21 13:42:36 -04:00
parent 2473d22d2e
commit 108e919dcb
8 changed files with 243 additions and 28 deletions
+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.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")
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}"
parts = [f"ai:{_clean(agent, 40) or 'agent'}"]
cleaned_action = _clean(action, 120)
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:
@@ -68,6 +108,10 @@ async def tag_vpgs(
"message": result.message,
"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]] = []
skipped = [v.as_dict() for v in result.vm.vpgs if not v.can_tag]
errors: list[str] = []
+59
View File
@@ -3,6 +3,7 @@
from __future__ import annotations
import asyncio
import urllib.parse
from typing import Any
from zerto_rewind_mcp.client import ZertoClient, ZertoError
@@ -94,3 +95,61 @@ async def wait_flr_ready(
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."
)
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."
)
+33 -10
View File
@@ -13,7 +13,12 @@ from zerto_rewind_mcp.checkpoints import checkpoint_id, checkpoint_tag, make_tag
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.recover import (
download_token_from,
resolve_flr_path,
session_id_from,
wait_flr_ready,
)
from zerto_rewind_mcp.util import pick
mcp = FastMCP("zerto_rewind_mcp")
@@ -110,12 +115,20 @@ 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.
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;
run this against the vSphere protected ZVM.
"""
@@ -125,7 +138,12 @@ async def zerto_create_tagged_checkpoint(
payload = result.as_dict()
payload["ok"] = False
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)
return _dump(out)
except ZertoError as exc:
@@ -182,13 +200,19 @@ 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)
return await zerto_create_tagged_checkpoint(
query=query, change_id=change_id, agent=agent, action=action
)
@mcp.tool(
@@ -261,8 +285,7 @@ async def zerto_recover_file(
"ok": False,
"needs_confirm": True,
"message": (
"Set confirmed=true after a human yes. "
"FLR mounts a disk on the recovery site."
"Set confirmed=true after a human yes. FLR mounts a disk on the recovery site."
),
}
)
@@ -279,8 +302,8 @@ async def zerto_recover_file(
)
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])
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"
@@ -292,6 +315,7 @@ async def zerto_recover_file(
"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}",
}
)
@@ -394,8 +418,7 @@ async def zerto_offsite_clone(
"ok": False,
"needs_confirm": True,
"message": (
"Set confirmed=true after a human yes. "
"Clone consumes recovery datastore."
"Set confirmed=true after a human yes. Clone consumes recovery datastore."
),
}
)