From 108e919dcb089578a395f7042eae506b1bf11682 Mon Sep 17 00:00:00 2001 From: "Claude (agent)" Date: Mon, 21 Sep 2026 13:42:36 -0400 Subject: [PATCH] fix(flr): partition-rooted FLR paths + agent/intent in checkpoint names (#2) --- CONTEXT.md | 2 +- README.md | 2 +- skills/zerto-rewind/SKILL.md | 16 +++++++- src/zerto_rewind_mcp/checkpoints.py | 54 +++++++++++++++++++++++--- src/zerto_rewind_mcp/recover.py | 59 +++++++++++++++++++++++++++++ src/zerto_rewind_mcp/server.py | 43 ++++++++++++++++----- tests/test_checkpoints.py | 46 +++++++++++++++++----- tests/test_recover.py | 49 ++++++++++++++++++++++++ 8 files changed, 243 insertions(+), 28 deletions(-) diff --git a/CONTEXT.md b/CONTEXT.md index dc5c3cc..6580ee9 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -5,7 +5,7 @@ PoC MCP that teaches an agent to discover Zerto protection, pin a tagged checkpo ## Language **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: | | vm= | change= | `. 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) **VPG**: diff --git a/README.md b/README.md index 05235e5..2b1e8c4 100644 --- a/README.md +++ b/README.md @@ -7,7 +7,7 @@ If the loop works, these tools are the delta to add to official ZVM MCP (`ZVM.MC ## 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. -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: | | vm= | change= | `. 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. diff --git a/skills/zerto-rewind/SKILL.md b/skills/zerto-rewind/SKILL.md index 835a544..75a926b 100644 --- a/skills/zerto-rewind/SKILL.md +++ b/skills/zerto-rewind/SKILL.md @@ -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: 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. 4. Then run the mutating call. @@ -51,4 +51,16 @@ Human must confirm. Pass `confirmed=true` only after they say yes. ## Tag -Default: `ai:::`. 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: | | vm= | change= | +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. diff --git a/src/zerto_rewind_mcp/checkpoints.py b/src/zerto_rewind_mcp/checkpoints.py index bf492b4..4b280f7 100644 --- a/src/zerto_rewind_mcp/checkpoints.py +++ b/src/zerto_rewind_mcp/checkpoints.py @@ -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] = [] diff --git a/src/zerto_rewind_mcp/recover.py b/src/zerto_rewind_mcp/recover.py index 0bf6920..e3e96d9 100644 --- a/src/zerto_rewind_mcp/recover.py +++ b/src/zerto_rewind_mcp/recover.py @@ -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." + ) diff --git a/src/zerto_rewind_mcp/server.py b/src/zerto_rewind_mcp/server.py index 1a2bec0..7be0f62 100644 --- a/src/zerto_rewind_mcp/server.py +++ b/src/zerto_rewind_mcp/server.py @@ -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:::. + + 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: | | vm= | change= | 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." ), } ) diff --git a/tests/test_checkpoints.py b/tests/test_checkpoints.py index 4eb5e2e..6663bf4 100644 --- a/tests/test_checkpoints.py +++ b/tests/test_checkpoints.py @@ -1,18 +1,46 @@ 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(): - when = datetime(2026, 9, 21, 15, 4, 5, tzinfo=UTC) - tag = make_tag("codex", "chg-99", when) - assert tag == "ai:codex:chg-99:20260921T150405Z" +def test_tag_describes_agent_and_action(): + tag = make_tag( + "claude", + "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(): - tag = make_tag("agent one", "path/with spaces", datetime(2026, 1, 1, tzinfo=UTC)) - assert " " not in tag - assert tag.startswith("ai:agent-one:path-with-spaces:") +def test_tag_without_action_or_vm(): + assert make_tag("codex", "chg-99", WHEN) == "ai:codex | change=chg-99 | 20260921T150405Z" + + +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(): diff --git a/tests/test_recover.py b/tests/test_recover.py index 4cf9b6b..361c501 100644 --- a/tests/test_recover.py +++ b/tests/test_recover.py @@ -28,3 +28,52 @@ def test_flr_list_status(): ] row = flr_row(payload) 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