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
This commit is contained in:
@@ -8,7 +8,7 @@ If the loop works, these tools are the delta to add to official ZVM MCP (`ZVM.MC
|
|||||||
|
|
||||||
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. The name records which agent and what it is doing: `ai:<agent> | <action> | vm=<vm> | change=<id> | <utc>`.
|
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.
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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}")
|
||||||
|
|||||||
@@ -153,3 +153,47 @@ async def resolve_flr_path(client: ZertoClient, session_id: str, guest_path: str
|
|||||||
f"{guest_path!r} not found in the FLR mount. Looked under {tried}. "
|
f"{guest_path!r} not found in the FLR mount. Looked under {tried}. "
|
||||||
"The file may not have replicated into that checkpoint yet."
|
"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
|
||||||
|
|||||||
@@ -15,8 +15,12 @@ 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 (
|
from zerto_rewind_mcp.recover import (
|
||||||
download_token_from,
|
download_token_from,
|
||||||
|
is_live_session,
|
||||||
resolve_flr_path,
|
resolve_flr_path,
|
||||||
session_id_from,
|
session_id_from,
|
||||||
|
session_id_of,
|
||||||
|
session_rows,
|
||||||
|
session_summary,
|
||||||
wait_flr_ready,
|
wait_flr_ready,
|
||||||
)
|
)
|
||||||
from zerto_rewind_mcp.util import pick
|
from zerto_rewind_mcp.util import pick
|
||||||
@@ -256,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={
|
||||||
@@ -292,7 +341,9 @@ async def zerto_recover_file(
|
|||||||
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,
|
||||||
@@ -309,8 +360,7 @@ async def zerto_recover_file(
|
|||||||
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),
|
||||||
@@ -318,15 +368,94 @@ async def zerto_recover_file(
|
|||||||
"flr_path": flr_path,
|
"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(
|
||||||
|
|||||||
@@ -77,3 +77,85 @@ def test_resolve_flr_path_picks_browsable_partition():
|
|||||||
assert got == "Volume2-Ext4/home/justin/app-config.yaml"
|
assert got == "Volume2-Ext4/home/justin/app-config.yaml"
|
||||||
# must not try the unrestorable partition
|
# must not try the unrestorable partition
|
||||||
assert "Volume1-Unknown/home/justin" not in client.seen
|
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