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

Merged
claude merged 2 commits from fix/flr-partition-path into main 2026-09-21 13:42:36 -04:00
8 changed files with 243 additions and 28 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**:
+1 -1
View File
@@ -7,7 +7,7 @@ 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`.
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.
+14 -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.
@@ -51,4 +51,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] = []
+59
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,61 @@ 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."
)
+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.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,
resolve_flr_path,
session_id_from,
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 +115,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 +138,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 +200,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(
@@ -261,8 +285,7 @@ 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."
), ),
} }
) )
@@ -279,8 +302,8 @@ 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"
@@ -292,6 +315,7 @@ async def zerto_recover_file(
"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}",
} }
) )
@@ -394,8 +418,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():
+49
View File
@@ -28,3 +28,52 @@ 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