Compare commits

..
2 Commits
Author SHA1 Message Date
justinandClaude Opus 5 0263f3797a docs(hooks): correct the cloud-protected denial claim
hooks/README.md said the win2019-1 denial happened because Zerto cannot
insert a tagged checkpoint on an AWS-protected VPG. That is not what
happened, and it is not true on 10.9.10.

The insert succeeded. The checkpoint is in the journal as cp 56, stamped
about a minute after the hook had already reported failure. What actually
happened is that wait_for_tag gave up after its hardcoded 45s, which is
generous on a VPG that checkpoints every 5s and far too short on one that
checkpoints every 630s.

So the documented "verified" example was a false denial. The deny path
works; that particular denial was wrong. Recorded as a known issue,
because a guard that silently refuses legitimate work on every
cloud-protected VM while looking correct is a worse failure than the one
it exists to prevent.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_016yVfC5nvZowoLFnEGWhLGn
2026-09-22 19:39:16 -04:00
justinandClaude Opus 5 a51b512bab feat(hooks): enforce the guard from a PreToolUse hook
The catalog and zerto_check_tool are advice. An MCP server cannot see or
block another server's tool calls, so a model that skips the guard is not
stopped by anything, and the change lands with no checkpoint behind it.

A PreToolUse hook runs in the host, where the tool call actually pauses.
That turns capture-before-execute from a convention into something the
host enforces.

  read-only                      no decision, runs
  mutating, checkpoint confirmed no decision, plus additionalContext
                                 telling the model which checkpoint to
                                 recover from
  mutating, checkpoint failed    DENY, the call never happens
  mutating, VM unknown to Zerto  prompt, nothing to rewind to
  mutating, no VM in the args    prompt, the catalog's vm_arg missed
  unlisted tool                  prompt

Decisions go out as JSON rather than exit 2. Exit 2 blocks
unconditionally but discards the JSON, and with it the reason, so the
model would be refused without being told why.

A lookup miss is not a refusal. Asking Zerto for a plain hostname by
vmIdentifier returns HTTP 400, and an early version reported that as
"could not tag a checkpoint", which denied changes to machines Zerto had
simply never heard of. Those are now separated: unknown VM prompts, a VM
Zerto knows but will not tag denies.

Timeouts are budgeted for the slow path. A successful tag took about 4s,
but a refusal took 63s, because wait_for_tag spends 45s waiting for a
checkpoint that will never arrive on an AWS or Azure protected VPG. If
the host's timeout fires first it cancels the hook and discards its
output, and the call proceeds unguarded, so the hook's own budget (150s)
stays under the configured one (180s): better to deny than be cancelled.

Verified against a live ZVM 10.9.10. jp-ubuntu tagged checkpoint 7180 and
was allowed; win2019-1, whose VPG is protected at an AWS site where
tagged checkpoints are unsupported, was denied. Both held under
permission_mode bypassPermissions, which is when an agent is most likely
running unattended.

Broad excepts in hooks/ are deliberate and scoped in pyproject: a hook
that raises breaks the tool call it exists to protect.

pytest 60 passed (12 new).

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_016yVfC5nvZowoLFnEGWhLGn
2026-09-22 19:39:16 -04:00
19 changed files with 49 additions and 1184 deletions
-10
View File
@@ -16,13 +16,3 @@ catalog.json
recovered/
.claude/*
!.claude/gitea-ship.json
# demo harness: live credentials and generated media never go in git
demo/demo_win.json
demo/vo/*.wav
demo/vo/narration.json
demo/vo/durations*.json
demo/*.cast
demo/*.gif
demo/*.mp4
demo/marks.jsonl
-75
View File
@@ -1,75 +0,0 @@
# Demo recording harness
Records the rewind loop running against a live ZVM as a narrated MP4. Terminal
only: no screen capture, no video editor.
```
tmux ──▶ asciinema ──▶ agg ──▶ ffmpeg ──▶ mp4
│ ▲
└── left pane: driver, right pane: journal │
│
xAI /v1/tts ──▶ wav per beat ──────────┘ (mux_vo.py)
```
## Pieces
| file | what it does |
|---|---|
| `windows_driver.py` | Windows demo. Intro over the diagram, then the live loop via WinRM. |
| `driver.py` | The Linux equivalent, over SSH. |
| `diagram.py` | Architecture diagram, revealed in four chunks against intro beats i1..i4. |
| `journal.py` | Right-hand pane. Polls a VPG's checkpoints so the tag appears on camera. |
| `record_win.sh` / `record.sh` | Drive tmux + asciinema, then agg and ffmpeg. |
| `narration.md` | The script. One `[beat]` per block; the source of truth. |
| `vo/build_narration.py` | Parses `narration.md`, synthesises a wav per beat, records durations. |
| `mux_vo.py` | Aligns the wavs to the recorded beat marks and muxes the audio. |
## Running it
```bash
cp demo/demo_win.example.json demo/demo_win.json # then fill in the guest creds
export XAI_KEY_FILE=~/xai-api.key XAI_VOICE_ID=<id from /v1/custom-voices>
python3 demo/vo/build_narration.py 1.0 # synthesise, writes durations.json
demo/record_win.sh win1 # record; writes marks.jsonl
python3 demo/mux_vo.py win1 # -> win1_narrated.mp4
```
`demo_win.json` holds live guest credentials and is gitignored. So are the
generated `.wav`, `.cast`, `.gif` and `.mp4` files.
## How the audio stays in sync
The driver writes `marks.jsonl` as it runs: one line per beat with the real
elapsed time it started. `mux_vo.py` delays each wav to its recorded mark, so
sync survives a slow API call or an FLR mount that takes longer than usual.
Nothing is predicted.
Two things this depends on:
- **`agg --idle-time-limit` must be larger than the longest pause** (the scripts
pass 3600). The default is 5 seconds, which compresses idle time, and that
silently breaks the mapping between wall clock and video time.
- **Each beat holds for its narration length.** `hold()` sleeps out whatever is
left after the work finishes, so a line is never cut off mid-sentence.
Check alignment after a mux: the FLR wait should measure near silence.
```bash
ffmpeg -v error -i out.mp4 -vn -ac 1 /tmp/a.wav
ffmpeg -hide_banner -ss 150 -t 6 -i /tmp/a.wav -af volumedetect -f null /dev/null 2>&1 | grep mean_volume
```
Speech sits around -22 dB; a correctly aligned gap reads about -91 dB.
## Narration gotchas
- **Do not map acronyms to run-together phonetics.** The `replace` map takes
`{"phrase": "pronunciation"}`, and `{"VM": "vee em"}` gets spoken as one word,
"vem". Either leave the acronym alone or expand it: `{"VM": "virtual machine"}`.
- **Write for speech, not for the page.** Short declaratives and fragment stacks
read well and sound robotic out loud. Commas and full stops are what the engine
uses for pacing, so clauses joined with commas breathe; a wall of four-word
sentences marches.
- `volumedetect` reports `n_samples: 0` when pointed at a file whose first stream
is video. Extract the audio first, then measure, or you will think a working
track is silent.
-12
View File
@@ -1,12 +0,0 @@
{
"vpg_id": "1816d7a4-3316-44dc-bafd-6d84f63d4ba9",
"vpg_name": "local",
"vm_query": "ad1(1)",
"vm_identifier": "f7f0835d-a5e4-46fa-b710-4bc32076e820.vm-2033",
"change_id": "win-demo-rewind",
"guest_path": "C:\\demo\\app-config.yaml",
"bad_content": "upstream: http://0.0.0.0:1 # BROKEN-BY-AGENT",
"host": "192.0.2.10",
"user": "EXAMPLE\\administrator",
"password": "change-me"
}
-44
View File
@@ -1,44 +0,0 @@
"""Architecture diagram, revealed in four chunks aligned to intro beats i1..i4."""
C = "\033[36m"; B = "\033[1m"; D = "\033[2m"; G = "\033[32m"; Y = "\033[33m"; R = "\033[0m"
CHUNKS = [
[
f" {B}{C}┌────────────────────────────────┐{R}",
f" {B}{C}│{R} Developer / process owner {B}{C}│{R} {D}\"change this config\"{R}",
f" {B}{C}└───────────────┬────────────────┘{R}",
f" {B}{C} │{R}",
f" {B}{C} ▼{R}",
f" {B}{C}┌────────────────────────────────┐{R}",
f" {B}{C}│{R} Business application on a VM {B}{C}│{R}",
f" {B}{C}└───────────────┬────────────────┘{R}",
],
[
f" {D} │ protected by{R}",
f" {D} ▼{R}",
f" {B}{G}┌────────────────────────────────┐{R} {D}┌───────────────────────────┐{R}",
f" {B}{G}│{R} Zerto continuous protection {B}{G}│{R} vs {D}│ Nightly backup │{R}",
f" {B}{G}│{R} {G}recovery point every few secs{R} {B}{G}│{R} {D}│ recovery point: hours old │{R}",
f" {B}{G}└────────────────────────────────┘{R} {D}└───────────────────────────┘{R}",
],
[
"",
f" {B}{Y}┌────────────────────────────────┐{R}",
f" {B}{Y}│{R} AI agent {B}{Y}│{R} 1. ask Zerto: is this protected?",
f" {B}{Y}│{R} authorised to make the change {B}{Y}│{R} 2. tag a checkpoint in the journal",
f" {B}{Y}└────────────────────────────────┘{R} 3. make the change, and test it",
],
[
"",
f" 4. {B}broken, and the agent cannot fix it{R}",
f" recover the file, or the whole VM,",
f" {G}from the checkpoint taken at step 2{R}",
],
]
def frame(upto: int) -> str:
"""Everything revealed through chunk `upto` (1-based)."""
out = ["\033[H\033[2J", f" {B}{C}ZERTO AI REWIND{R} {D}how it fits together{R}", ""]
for chunk in CHUNKS[:upto]:
out.extend(chunk)
return "\n".join(out)
-254
View File
@@ -1,254 +0,0 @@
"""Left pane: drive zerto_rewind_mcp over real MCP stdio. Full rewind loop."""
from __future__ import annotations
import asyncio, json, os, shlex, subprocess, sys, time
from pathlib import Path
from mcp import ClientSession, StdioServerParameters
from mcp.client.stdio import stdio_client
SP = Path(__file__).parent
CFG = json.loads((SP / "demo.json").read_text())
REPO = "/home/justin/github/zerto-ai-rewind"
RST, BOLD, DIM = "\033[0m", "\033[1m", "\033[2m"
CYAN, GREEN, RED, YEL, MAG = "\033[36m", "\033[32m", "\033[31m", "\033[33m", "\033[35m"
SPEED = float(os.environ.get("DEMO_SPEED", "1.0"))
def recovered_bytes(rec: dict) -> bytes:
"""zerto_recover_file returns the content now, not a path on the MCP host.
The server may be on another machine, so a path there is of no use here.
"""
import base64
if rec.get("encoding") == "base64":
return base64.b64decode(rec["content"])
return rec["content"].encode()
def w(s=""):
sys.stdout.write(s + "\n"); sys.stdout.flush()
def type_out(s, delay=0.012, color=""):
sys.stdout.write(color)
for ch in s:
sys.stdout.write(ch); sys.stdout.flush()
time.sleep(delay * SPEED)
sys.stdout.write(RST + "\n"); sys.stdout.flush()
def beat(n=1.2):
time.sleep(n * SPEED)
STEP = [0]
def step(title):
STEP[0] += 1
w()
w(f"{BOLD}{CYAN}{'='*66}{RST}")
type_out(f" STEP {STEP[0]} {title}", 0.008, BOLD + CYAN)
w(f"{BOLD}{CYAN}{'='*66}{RST}")
beat(0.5)
def note(s):
type_out(f" {s}", 0.010, DIM)
def call_banner(tool, args):
w(f"{MAG} -> MCP call{RST} {BOLD}{tool}{RST}")
for k, v in args.items():
w(f"{DIM} {k} = {v}{RST}")
def show(payload, keep=None, limit=22):
if keep:
payload = {k: payload[k] for k in keep if k in payload}
text = json.dumps(payload, indent=2, default=str)
lines = text.splitlines()
for line in lines[:limit]:
w(f"{DIM} |{RST} {line}")
if len(lines) > limit:
w(f"{DIM} | ... {len(lines)-limit} more lines{RST}")
def ssh_base():
if CFG.get("password"):
return ["sshpass", "-p", CFG["password"], "ssh", "-o", "StrictHostKeyChecking=no",
"-o", "UserKnownHostsFile=/dev/null", "-o", "LogLevel=ERROR"]
base = ["ssh", "-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null",
"-o", "LogLevel=ERROR"]
if CFG.get("key"):
base += ["-i", CFG["key"]]
return base
def guest(cmd, show_cmd=True):
target = f"{CFG['user']}@{CFG['host']}"
full = ssh_base() + [target, cmd]
if show_cmd:
w(f"{YEL} $ ssh {target} {shlex.quote(cmd)}{RST}")
r = subprocess.run(full, capture_output=True, text=True, timeout=60)
out = (r.stdout + r.stderr).rstrip()
for line in out.splitlines():
w(f"{DIM} |{RST} {line}")
return r.returncode, out
async def main():
w(f"\n{BOLD}{CYAN} ZERTO AI REWIND{RST} {DIM}live end-to-end{RST}\n")
type_out(" An agent is about to change a file on a protected VM.", 0.014)
type_out(" Zerto already journals that VM. Make the agent use it.", 0.014)
beat(1.5)
params = StdioServerParameters(
command=f"{REPO}/.venv/bin/zerto-rewind-mcp",
env={**os.environ, "ZERTO_REWIND_CONFIG": f"{REPO}/config.json"},
)
devnull = open(os.devnull, 'w')
async with stdio_client(params, errlog=devnull) as (r, s):
async with ClientSession(r, s) as sess:
step("Connect to the rewind MCP")
await sess.initialize()
tools = await sess.list_tools()
note(f"{len(tools.tools)} tools from zerto_rewind_mcp")
for t in tools.tools:
w(f"{DIM} - {t.name}{RST}")
beat(2.0)
async def _ticker(label):
t0 = time.time()
try:
while True:
await asyncio.sleep(5)
el = int(time.time() - t0)
sys.stdout.write(f"\r{DIM} {label} ... {el}s{RST}")
sys.stdout.flush()
except asyncio.CancelledError:
el = int(time.time() - t0)
if el >= 5:
sys.stdout.write(f"\r{DIM} {label} ... {el}s done{RST}\n")
else:
sys.stdout.write("\r" + " " * 50 + "\r")
sys.stdout.flush()
raise
async def call(name, **args):
call_banner(name, args)
tick = asyncio.create_task(_ticker("working"))
try:
res = await sess.call_tool(name, args)
finally:
tick.cancel()
try:
await tick
except asyncio.CancelledError:
pass
return json.loads(res.content[0].text)
# ---- 1. the good state on the guest
step("The file the agent is about to break")
guest(f"cat {CFG['guest_path']}")
beat(1.8)
# ---- 2. find protection
step("zerto_find_protection - is this VM protected?")
note("Reads skip the guard. This is the read.")
found = await call("zerto_find_protection", query=CFG["vm_query"])
show(found)
if not found.get("ok"):
w(f"{RED} find_protection not ok - stopping{RST}"); return 1
vm = found["vm"]
beat(2.5)
# ---- 3. guard
step("zerto_guard_before_mutate - pin the journal FIRST")
note("Capture-before-execute. No tag, no change.")
note("The checkpoint name records which agent, and what it is about to do.")
note("Watch the journal pane on the right.")
guard = await call("zerto_guard_before_mutate",
query=CFG["vm_query"], change_id=CFG["change_id"], agent="claude",
action=f"edit {CFG['guest_path']}")
show(guard, keep=["ok", "tag", "tagged", "skipped", "errors", "message"])
if not guard.get("ok"):
w(f"{RED} GUARD FAILED - refusing the change{RST}"); return 1
tag = guard["tag"]
tagged = guard["tagged"][0]
t_tag = time.strftime("%H:%M:%S")
w(f"{GREEN}{BOLD} tagged checkpoint {tagged['checkpoint_id']} on {tagged['vpg_name']}{RST}")
w(f"{GREEN} inserted at {t_tag}{RST}")
w(f"{GREEN} name: {tag}{RST}")
w(f"{DIM} this checkpoint is the rewind point. remember cp {tagged['checkpoint_id']}.{RST}")
beat(4.0)
# ---- 4. the bad change
step("Now the agent makes the bad change")
guest(f"printf '%s\\n' {shlex.quote(CFG['bad_content'])} > {CFG['guest_path']}")
guest(f"cat {CFG['guest_path']}")
t_bad = time.strftime("%H:%M:%S")
w(f"{RED}{BOLD} the good content is gone from the guest ({t_bad}){RST}")
note("Git never had this file. Last night's backup is hours stale.")
beat(2.5)
# ---- 5. refusal without a human
step("zerto_recover_file without a human yes")
deny = await call("zerto_recover_file",
vpg_identifier=tagged["vpg_identifier"],
vm_identifier=vm["vm_identifier"],
checkpoint_identifier=tagged["checkpoint_id"],
guest_path=CFG["guest_path"], confirmed=False)
show(deny)
w(f"{YEL} refused. FLR mounts a disk; a human says yes.{RST}")
beat(2.5)
# ---- 6. recover
rec_cp = tagged["checkpoint_id"]
step("Human says yes - FLR from the pre-mutation checkpoint")
w(f"{BOLD} provenance{RST}")
w(f" {GREEN}cp {tagged['checkpoint_id']} inserted {t_tag}{RST} {DIM}<- guard, BEFORE the change{RST}")
w(f" {RED}bad write {t_bad}{RST} {DIM}<- the mutation, AFTER{RST}")
assert rec_cp == tagged["checkpoint_id"], "recovery must use the guard checkpoint"
w(f"{DIM} recovering from cp {rec_cp}, not the newest checkpoint.{RST}")
note(f"name: {tag}")
rec = await call("zerto_recover_file",
vpg_identifier=tagged["vpg_identifier"],
vm_identifier=vm["vm_identifier"],
checkpoint_identifier=tagged["checkpoint_id"],
guest_path=CFG["guest_path"], confirmed=True)
show(rec)
if not rec.get("ok"):
w(f"{RED} FLR failed{RST}"); return 1
beat(1.5)
w(f"{YEL} recovered {rec['bytes']} bytes, sha256 {rec['sha256'][:16]}...{RST}")
body = recovered_bytes(rec).decode()
for line in body.splitlines():
w(f"{GREEN}{BOLD} | {line}{RST}")
beat(2.5)
# ---- 7. put it back
step("Copy it back to the guest")
note("This part is scp, not Zerto. Zerto got the bytes back.")
target = f"{CFG['user']}@{CFG['host']}"
local_copy = SP / rec["name"]
local_copy.write_bytes(recovered_bytes(rec))
scp = (["sshpass", "-p", CFG["password"]] if CFG.get("password") else []) + \
["scp", "-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null",
"-o", "LogLevel=ERROR"] + (["-i", CFG["key"]] if CFG.get("key") else []) + \
[str(local_copy), f"{target}:{CFG['guest_path']}"]
w(f"{YEL} $ scp {rec['name']} {target}:{CFG['guest_path']}{RST}")
subprocess.run(scp, capture_output=True, text=True, timeout=60)
guest(f"cat {CFG['guest_path']}")
beat(1.5)
w()
w(f"{BOLD}{GREEN}{'='*66}{RST}")
type_out(" REWOUND. RPO was the journal, not last night.", 0.016, BOLD + GREEN)
w(f"{BOLD}{GREEN}{'='*66}{RST}")
beat(4.0)
return 0
sys.exit(asyncio.run(main()))
-78
View File
@@ -1,78 +0,0 @@
"""Right pane: live view of a VPG journal, tagged checkpoints highlighted."""
from __future__ import annotations
import asyncio, os, sys, time
os.environ.setdefault("ZERTO_REWIND_CONFIG", "/home/justin/github/zerto-ai-rewind/config.json")
from zerto_rewind_mcp.config import load_config
from zerto_rewind_mcp.client import ZertoClient
from zerto_rewind_mcp.checkpoints import checkpoint_tag, checkpoint_id
from zerto_rewind_mcp.util import pick
VPG_ID = sys.argv[1]
VPG_NAME = sys.argv[2] if len(sys.argv) > 2 else VPG_ID
ROWS = int(os.environ.get("JOURNAL_ROWS", "10"))
DIM, RST, BOLD = "\033[2m", "\033[0m", "\033[1m"
CYAN, GREEN, YEL = "\033[36m", "\033[32m", "\033[33m"
def _wrap_tag(tag, width=52, maxlines=3):
tag = tag.replace("; Used for File Level Restore", " [FLR]")
out = []
while tag and len(out) < maxlines:
out.append(tag[:width]); tag = tag[width:]
return out
def ts_of(row):
return str(pick(row, "TimeStamp", "Timestamp", "timestamp") or "")
def render(rows, spin, tagged):
out = ["\033[H\033[2J"]
out.append(f"{BOLD}{CYAN} ZERTO JOURNAL {RST}{BOLD}{VPG_NAME}{RST}\n")
out.append(f"{DIM} vpg {VPG_ID}{RST}\n")
out.append(f"{DIM} {len(rows)} checkpoints{RST} {YEL}{len(tagged)} ai-tagged{RST} {DIM}{spin}{RST}\n\n")
out.append(f"{BOLD} AI TAGS{RST}\n")
if not tagged:
out.append(f"{DIM} (none yet){RST}\n")
for row in tagged[-4:]:
t = ts_of(row).replace("T", " ").replace(".000Z", "")[11:]
out.append(f"{GREEN}{BOLD} * {t} cp {checkpoint_id(row)}{RST}\n")
for seg in _wrap_tag(checkpoint_tag(row)):
out.append(f"{GREEN} {seg}{RST}\n")
out.append(f"\n{BOLD} NEWEST{RST}\n")
for row in rows[-ROWS:]:
tag = checkpoint_tag(row)
t = ts_of(row).replace("T", " ").replace(".000Z", "")
cid = checkpoint_id(row)
t = t[11:]
if tag.startswith("ai:"):
out.append(f"{GREEN}{BOLD} * {t} cp {cid} <-- new{RST}\n")
elif tag:
out.append(f"{DIM} {t} cp {cid} {tag[:20]}{RST}\n")
else:
out.append(f"{DIM} {t} cp {cid}{RST}\n")
sys.stdout.write("".join(out))
sys.stdout.flush()
async def main():
s = load_config()
c = ZertoClient(base_url=s["zerto_url"], username=s["username"], password=s["password"],
client_id=s.get("client_id", "zerto-client"), verify_tls=bool(s.get("verify_tls")))
spinner = "|/-\\"
i = 0
while True:
try:
rows = await c.list_checkpoints(VPG_ID)
tagged = [r for r in rows if checkpoint_tag(r).startswith("ai:")]
render(rows, spinner[i % 4], tagged)
except Exception as exc:
sys.stdout.write(f"\n journal poll error: {exc}\n")
sys.stdout.flush()
i += 1
await asyncio.sleep(3)
asyncio.run(main())
-39
View File
@@ -1,39 +0,0 @@
"""Build the narration track from recorded beat marks and mux it onto the video."""
import json, subprocess, sys
from pathlib import Path
SP = Path(__file__).parent
VO = SP / "vo"
take = sys.argv[1] if len(sys.argv) > 1 else "win1"
video = SP / f"{take}.mp4"
out = SP / f"{take}_narrated.mp4"
marks = [json.loads(l) for l in (SP / "marks.jsonl").read_text().splitlines() if l.strip()]
dur = json.loads((VO / "durations.json").read_text())
vlen = float(subprocess.run(["ffprobe", "-v", "error", "-show_entries", "format=duration",
"-of", "default=nw=1:nk=1", str(video)],
capture_output=True, text=True).stdout.strip())
inputs, filters, labels = [], [], []
for i, m in enumerate(marks):
wav = VO / f"{m['key']}.wav"
if not wav.exists():
print(f" missing {wav.name}, skipping"); continue
delay_ms = int(m["t"] * 1000)
inputs += ["-i", str(wav)]
filters.append(f"[{i}:a]adelay={delay_ms}|{delay_ms},apad[a{i}]")
labels.append(f"[a{i}]")
end = m["t"] + dur.get(m["key"], 0)
flag = "" if end <= vlen + 0.5 else " <-- OVERRUNS VIDEO"
print(f" {m['key']:9} start {m['t']:7.2f}s len {dur.get(m['key'],0):5.2f}s end {end:7.2f}s{flag}")
mix = "".join(labels) + f"amix=inputs={len(labels)}:normalize=0,atrim=0:{vlen},asetpts=N/SR/TB[a]"
cmd = (["ffmpeg", "-y", "-loglevel", "error"] + inputs + ["-i", str(video),
"-filter_complex", ";".join(filters) + ";" + mix,
"-map", f"{len(labels)}:v", "-map", "[a]",
"-c:v", "copy", "-c:a", "aac", "-b:a", "160k", "-shortest", str(out)])
r = subprocess.run(cmd, capture_output=True, text=True)
if r.returncode != 0:
print("ffmpeg failed:\n", r.stderr[-1500:]); sys.exit(1)
print(f"\n video {vlen:.1f}s | narration ends {max(m['t']+dur.get(m['key'],0) for m in marks):.1f}s")
print(f" -> {out}")
-52
View File
@@ -1,52 +0,0 @@
# Zerto AI Rewind - demo narration (v3, rewritten against the jp-voice profile)
# Edit any line, then tell Claude to re-read this file.
# Blank line separates beats. Lines starting with # are ignored.
## PART 1 - INTRO (architecture diagram on screen, builds as each beat lands)
[i1] Hello, and welcome to the Zerto AI Rewind demo.
[i2] AI agents are acting on production systems faster than ever, and a lot faster
than backup was ever designed to keep up with. Zerto's continuous data
protection is a good fit for that, because the journal is always running.
[i3] In this demo an AI agent makes a change to a production system. But before it
does, it checks whether that VM is protected by Zerto, and if it is, it inserts
a tagged checkpoint into the journal first.
[i4] That gives the agent something to fall back on. If the change breaks the
application, and the agent can't fix it on its own, it can use file level or
full system recovery to get back to the moment before it touched anything.
## PART 2 - THE DEMO (audience: a customer, not an engineer)
[d1] This is the agent connecting to Zerto. Everything from here runs against a live
Zerto environment.
[d2] The machine is a production Windows server, and Zerto is already protecting it.
The agent has been asked to change a configuration file on that server.
[d3] Before it changes anything, it asks Zerto a simple question. Is this machine
protected, and which protection group is it in?
[d4] It is, so Zerto puts a tagged checkpoint into the journal, and the agent waits
until Zerto confirms that checkpoint is really there before it goes any further.
[d5] Now the change goes in, and it breaks the file. The version that was there is
gone from the server. Backup has a copy from last night, so that's hours old
already. Zerto has one from seconds before the change.
[d6] The agent can't just recover on its own. Recovery has to be approved by a
person, so at this point it stops and asks.
[d7] With approval, Zerto goes back to the checkpoint from just before the change,
and hands back the original file exactly as it was.
[d8] Putting that file back onto the server is an ordinary copy. Zerto's job was
keeping the data in the first place.
[d9] This was a single file, so file level recovery was enough. That same checkpoint
would let you bring back the entire machine if the damage were bigger.
[d10] And that's the point. The recovery point is the journal, seconds before the
change, instead of last night's backup window.
-52
View File
@@ -1,52 +0,0 @@
#!/usr/bin/env bash
# Record the split-pane rewind demo to an MP4.
set -uo pipefail
SP="$(cd "$(dirname "$0")" && pwd)"
REPO=/home/justin/github/zerto-ai-rewind
PY="$REPO/.venv/bin/python"
VPG_ID="$(python3 -c "import json;print(json.load(open('$SP/demo.json'))['vpg_id'])")"
VPG_NAME="$(python3 -c "import json;print(json.load(open('$SP/demo.json'))['vpg_name'])")"
TAKE="${1:-take1}"
CAST="$SP/$TAKE.cast"; GIF="$SP/$TAKE.gif"; MP4="$SP/$TAKE.mp4"
tmux -L demo kill-server 2>/dev/null
tmux -L rec kill-server 2>/dev/null
sleep 1
# demo session: driver left, journal right
tmux -L demo new-session -d -s demo -x 200 -y 50 -c "$REPO" \
"$PY $SP/driver.py; echo; echo ' [take complete]'; sleep 3; tmux -L demo kill-server"
tmux -L demo set-option -g status off
tmux -L demo split-window -h -l 60 -t demo -c "$REPO" \
"JOURNAL_ROWS=10 $PY $SP/journal.py $VPG_ID $VPG_NAME"
tmux -L demo select-pane -t 0
# recorder session gives asciinema a pty
tmux -L rec new-session -d -s rec -x 200 -y 50 \
"$SP/recvenv/bin/asciinema rec '$CAST' --overwrite --quiet -c 'tmux -L demo attach -t demo'"
echo "recording -> $CAST"
for i in $(seq 1 600); do
tmux -L rec has-session -t rec 2>/dev/null || break
sleep 2
done
tmux -L demo kill-server 2>/dev/null
tmux -L rec kill-server 2>/dev/null
sleep 1
[ -s "$CAST" ] || { echo "NO CAST PRODUCED"; exit 1; }
echo "cast: $(du -h "$CAST" | cut -f1) duration: $(python3 -c "
import json,sys
last=0
for l in open('$CAST'):
l=l.strip()
if l.startswith('['):
last=json.loads(l)[0]
print(f'{last:.0f}s')")"
"$SP/bin/agg" --font-size 14 --fps-cap 10 --idle-time-limit 2 --theme asciinema "$CAST" "$GIF"
ffmpeg -y -loglevel error -i "$GIF" \
-movflags +faststart -pix_fmt yuv420p -c:v libx264 -crf 20 \
-vf "scale=trunc(iw/2)*2:trunc(ih/2)*2" "$MP4"
echo "MP4: $MP4 ($(du -h "$MP4" | cut -f1))"
ffprobe -v error -show_entries format=duration:stream=width,height -of default=nw=1 "$MP4"
-42
View File
@@ -1,42 +0,0 @@
#!/usr/bin/env bash
# Record the Windows rewind demo. NO idle compression: the video timeline must
# equal wall clock so the narration can be aligned to the recorded beat marks.
set -uo pipefail
SP="$(cd "$(dirname "$0")" && pwd)"
REPO=/home/justin/github/zerto-ai-rewind
PY="$REPO/.venv/bin/python"
VPG_ID=$(python3 -c "import json;print(json.load(open('$SP/demo_win.json'))['vpg_id'])")
TAKE="${1:-win1}"
CAST="$SP/$TAKE.cast"; GIF="$SP/$TAKE.gif"; MP4="$SP/$TAKE.mp4"
tmux -L demo kill-server 2>/dev/null; tmux -L rec kill-server 2>/dev/null; sleep 1
rm -f "$SP/marks.jsonl"
tmux -L demo new-session -d -s demo -x 200 -y 50 -c "$REPO" \
"$PY $SP/windows_driver.py; echo; echo ' [take complete]'; sleep 2; tmux -L demo kill-server"
tmux -L demo set-option -g status off
tmux -L demo split-window -h -l 60 -t demo -c "$REPO" \
"JOURNAL_ROWS=10 $PY $SP/journal.py $VPG_ID ad1"
tmux -L demo select-pane -t 0
tmux -L rec new-session -d -s rec -x 200 -y 50 \
"$SP/recvenv/bin/asciinema rec '$CAST' --overwrite --quiet -c 'tmux -L demo attach -t demo'"
echo "recording -> $CAST"
for i in $(seq 1 450); do tmux -L rec has-session -t rec 2>/dev/null || break; sleep 2; done
tmux -L demo kill-server 2>/dev/null; tmux -L rec kill-server 2>/dev/null; sleep 1
[ -s "$CAST" ] || { echo "NO CAST"; exit 1; }
echo "cast duration: $(python3 -c "
import json
last=0
for l in open('$CAST'):
l=l.strip()
if l.startswith('['): last=json.loads(l)[0]
print(f'{last:.1f}s')")"
# idle-time-limit huge = no compression, so 1s of cast == 1s of video
"$SP/bin/agg" --font-size 14 --fps-cap 10 --idle-time-limit 3600 --theme asciinema "$CAST" "$GIF"
ffmpeg -y -loglevel error -i "$GIF" -movflags +faststart -pix_fmt yuv420p -c:v libx264 -crf 20 \
-vf "scale=trunc(iw/2)*2:trunc(ih/2)*2" "$MP4"
echo "MP4: $MP4 ($(du -h "$MP4" | cut -f1))"
ffprobe -v error -show_entries format=duration:stream=width,height -of default=nw=1 "$MP4"
-35
View File
@@ -1,35 +0,0 @@
"""Parse the approved markdown script into clips + durations."""
import json, os, re, subprocess, sys
from pathlib import Path
SP = Path(__file__).parent
SRC = Path(os.environ.get("NARRATION_MD", SP.parent / "narration.md"))
KEY = Path(os.environ.get("XAI_KEY_FILE", "~/xai-api.key")).expanduser().read_text().strip()
VOICE = os.environ.get("XAI_VOICE_ID", "") # a custom voice id from /v1/custom-voices
SPEED = float(sys.argv[1]) if len(sys.argv) > 1 else 1.0
beats = re.findall(r"\[([id]\d+)\]\s*(.*?)(?=\n\n|\Z)", SRC.read_text(), re.S)
lines = {k: " ".join(v.split()) for k, v in beats}
(SP / "narration.json").write_text(json.dumps(lines, indent=2) + "\n")
# Do NOT map an acronym to run-together phonetics: {"VM": "vee em"} is spoken
# as one word, "vem". Expand it instead, or leave it alone.
REPLACE = {"VM": "virtual machine"}
out = {}
for key, text in lines.items():
body = {"text": text, "voice_id": VOICE, "language": "en", "speed": SPEED,
"replace": REPLACE, "output_format": {"codec": "wav", "sample_rate": 24000}}
req = SP / "req.json"; req.write_text(json.dumps(body))
dest = SP / f"{key}.wav"
code = subprocess.run(["curl", "-sS", "-o", str(dest), "-w", "%{http_code}", "-X", "POST",
"https://api.x.ai/v1/tts", "-H", f"Authorization: Bearer {KEY}",
"-H", "Content-Type: application/json", "--data-binary", f"@{req}"],
capture_output=True, text=True).stdout.strip()
req.unlink()
if code != "200":
print(f" {key}: HTTP {code} FAILED"); sys.exit(1)
d = float(subprocess.run(["ffprobe", "-v", "error", "-show_entries", "format=duration",
"-of", "default=nw=1:nk=1", str(dest)], capture_output=True, text=True).stdout.strip())
out[key] = round(d, 2)
print(f" {key:5} {d:5.2f}s {len(text.split()):3} words")
(SP / "durations.json").write_text(json.dumps(out, indent=1))
print(f"\n {len(out)} clips, {sum(out.values()):.1f}s of narration")
-204
View File
@@ -1,204 +0,0 @@
"""Windows rewind demo. Beats hold for their narration length and log real start times."""
from __future__ import annotations
import json, os, sys, time
from pathlib import Path
import winrm
import diagram
from mcp import ClientSession, StdioServerParameters
from mcp.client.stdio import stdio_client
SP = Path(__file__).parent
CFG = json.loads((SP / "demo_win.json").read_text())
DUR = json.loads((SP / "vo" / "durations.json").read_text())
REPO = "/home/justin/github/zerto-ai-rewind"
MARKS = SP / "marks.jsonl"
PAD = 0.8 # breath between beats
RST, BOLD, DIM = "\033[0m", "\033[1m", "\033[2m"
CYAN, GREEN, RED, YEL, MAG = "\033[36m", "\033[32m", "\033[31m", "\033[33m", "\033[35m"
T0 = time.time()
STEP = [0]
def recovered_bytes(rec: dict) -> bytes:
"""zerto_recover_file returns the content now, not a path on the MCP host.
The server may be on another machine, so a path there is of no use here.
"""
import base64
if rec.get("encoding") == "base64":
return base64.b64decode(rec["content"])
return rec["content"].encode()
def w(s=""):
sys.stdout.write(s + "\n"); sys.stdout.flush()
def type_out(s, delay=0.018, color=""):
sys.stdout.write(color)
for ch in s:
sys.stdout.write(ch); sys.stdout.flush(); time.sleep(delay)
sys.stdout.write(RST + "\n"); sys.stdout.flush()
def begin(key, title=None):
"""Mark the beat start so the narration can be aligned to it later."""
t = time.time() - T0
with MARKS.open("a") as fh:
fh.write(json.dumps({"key": key, "t": round(t, 3)}) + "\n")
if title:
STEP[0] += 1
w(); w(f"{BOLD}{CYAN}{'=' * 74}{RST}")
type_out(f" STEP {STEP[0]} {title}", 0.010, BOLD + CYAN)
w(f"{BOLD}{CYAN}{'=' * 74}{RST}")
return t
def hold(start_t, key):
"""Keep this beat on screen until its narration has finished."""
want = DUR.get(key, 6.0) + PAD
spent = (time.time() - T0) - start_t
if spent < want:
time.sleep(want - spent)
def show(payload, keep=None, limit=18):
if keep:
payload = {k: payload[k] for k in keep if k in payload}
for line in json.dumps(payload, indent=2, default=str).splitlines()[:limit]:
w(f"{DIM} |{RST} {line}")
_sess = winrm.Session(f"http://{CFG['host']}:5985/wsman",
auth=(CFG["user"], CFG["password"]), transport="ntlm")
def guest_ps(ps, label):
w(f"{YEL} PS {CFG['host']}> {label}{RST}")
r = _sess.run_ps(ps)
out = r.std_out.decode(errors="replace").strip()
for line in out.splitlines():
w(f"{DIM} |{RST} {line}")
return out
async def main():
import asyncio
MARKS.unlink(missing_ok=True)
for n, key in enumerate(("i1", "i2", "i3", "i4"), start=1):
t = begin(key)
sys.stdout.write(diagram.frame(n)); sys.stdout.write("\n"); sys.stdout.flush()
hold(t, key)
params = StdioServerParameters(
command=f"{REPO}/.venv/bin/zerto-rewind-mcp",
env={**os.environ, "ZERTO_REWIND_CONFIG": f"{REPO}/config.json"})
devnull = open(os.devnull, "w")
async with stdio_client(params, errlog=devnull) as (r, s):
async with ClientSession(r, s) as sess:
t = begin("d1", "Connect to the rewind MCP")
await sess.initialize()
tools = await sess.list_tools()
w(f"{DIM} {len(tools.tools)} tools from zerto_rewind_mcp{RST}")
for x in tools.tools:
w(f"{DIM} - {x.name}{RST}")
hold(t, "d1")
async def call(name, **args):
w(f"{MAG} -> MCP call{RST} {BOLD}{name}{RST}")
for k, v in args.items():
w(f"{DIM} {k} = {v}{RST}")
res = await sess.call_tool(name, args)
return json.loads(res.content[0].text)
t = begin("d2", "The file the agent is about to break")
guest_ps(f"Get-Content '{CFG['guest_path']}'", f"Get-Content {CFG['guest_path']}")
hold(t, "d2")
t = begin("d3", "zerto_find_protection")
found = await call("zerto_find_protection", query=CFG["vm_query"])
show(found, keep=["ok", "outcome", "vm"])
if not found.get("ok"):
w(f"{RED} not ok, stopping{RST}"); return 1
hold(t, "d3")
t = begin("d4", "zerto_guard_before_mutate")
guard = await call("zerto_guard_before_mutate", query=CFG["vm_query"],
change_id=CFG["change_id"], agent="claude",
action=f"edit {CFG['guest_path']}")
show(guard, keep=["ok", "tag", "tagged"])
if not guard.get("ok"):
w(f"{RED} GUARD FAILED, refusing the change{RST}"); return 1
tg = guard["tagged"][0]
t_tag = time.strftime("%H:%M:%S")
w(f"{GREEN}{BOLD} checkpoint {tg['checkpoint_id']} inserted {t_tag} task={tg['task_state']}{RST}")
hold(t, "d4")
t = begin("d5", "Now the agent makes the bad change")
guest_ps(f"Set-Content -Path '{CFG['guest_path']}' -Value '{CFG['bad_content']}' -Encoding ASCII; "
f"Get-Content '{CFG['guest_path']}'", f"Set-Content {CFG['guest_path']} ...")
t_bad = time.strftime("%H:%M:%S")
w(f"{RED}{BOLD} the good config is gone from the guest ({t_bad}){RST}")
hold(t, "d5")
t = begin("d6", "zerto_recover_file without a human yes")
deny = await call("zerto_recover_file", vpg_identifier=CFG["vpg_id"],
vm_identifier=CFG["vm_identifier"],
checkpoint_identifier=tg["checkpoint_id"],
guest_path=CFG["guest_path"], confirmed=False)
show(deny)
w(f"{YEL} refused.{RST}")
hold(t, "d6")
t = begin("d7", "Human says yes: FLR from the pre-mutation checkpoint")
w(f"{BOLD} provenance{RST}")
w(f" {GREEN}cp {tg['checkpoint_id']} inserted {t_tag}{RST} {DIM}<- guard, BEFORE{RST}")
w(f" {RED}bad write {t_bad}{RST} {DIM}<- the mutation, AFTER{RST}")
rec = await call("zerto_recover_file", vpg_identifier=CFG["vpg_id"],
vm_identifier=CFG["vm_identifier"],
checkpoint_identifier=tg["checkpoint_id"],
guest_path=CFG["guest_path"], confirmed=True,
dest_dir=str(SP / "win_demo_out"))
show(rec, keep=["ok", "bytes", "flr_path", "unmount"])
if not rec.get("ok"):
w(f"{RED} FLR failed{RST}"); return 1
body = recovered_bytes(rec).decode()
w(f"{YEL} recovered file:{RST}")
for line in body.splitlines():
w(f"{GREEN}{BOLD} | {line}{RST}")
hold(t, "d7")
t = begin("d8", "Copy it back to the guest")
# Ship the recovered bytes verbatim. Set-Content rewrites line
# endings and appends one, which makes the restored file differ
# from what Zerto handed back. WriteAllBytes does not.
import base64
blob = recovered_bytes(rec)
b64 = base64.b64encode(blob).decode()
guest_ps(
f"[IO.File]::WriteAllBytes('{CFG['guest_path']}',"
f"[Convert]::FromBase64String('{b64}'));"
f"Get-Content '{CFG['guest_path']}';"
f"'md5: ' + (Get-FileHash '{CFG['guest_path']}' -Algorithm MD5).Hash",
f"WriteAllBytes {CFG['guest_path']} <{len(blob)} recovered bytes>")
hold(t, "d8")
t = begin("d9")
w(); w(f"{BOLD}{YEL} the same checkpoint, bigger blast radius{RST}")
type_out(" one file -> file level recovery (what you just saw)", 0.012, DIM)
type_out(" the whole VM -> offsite clone, failover test, or failover", 0.012, DIM)
hold(t, "d9")
t = begin("d10")
w(); w(f"{BOLD}{GREEN}{'=' * 74}{RST}")
type_out(" REWOUND. The recovery point was the journal.", 0.018, BOLD + GREEN)
w(f"{BOLD}{GREEN}{'=' * 74}{RST}")
hold(t, "d10")
return 0
import asyncio
sys.exit(asyncio.run(main()))
+23 -25
View File
@@ -37,19 +37,14 @@ venv that has this package installed.
The hook is synchronous: the host waits. That is the point, because the
checkpoint has to exist before the change does.
Budget for the slow path, not the fast one. How long the guard takes is set by
the VPG's checkpoint cadence, which is set in turn by its protected site:
Budget for the slow path, not the fast one. A successful tag took about 4s
against a healthy vSphere-protected VPG, but a **refusal took 63s**, because
`wait_for_tag` spends 45s before giving up.
| protected at | cadence | guard takes |
|---|---|---|
| vSphere | 5s | ~7s |
| Azure | 60s | ~40s |
| AWS | 630s | ~111s |
So:
The tag wait is derived from that cadence and capped at 300s, so:
- `timeout` in settings.json: 360 (seconds)
- `ZERTO_HOOK_GUARD_TIMEOUT`: 330 (seconds), kept under it
- `timeout` in settings.json: 180 (seconds)
- `ZERTO_HOOK_GUARD_TIMEOUT`: 150 (seconds), kept under it
If the host's timeout fires first it cancels the hook and **discards its
output**, and the tool call carries on through the normal permission flow. A
@@ -73,23 +68,26 @@ likely to be running unattended.
`~/.zerto-guard-hook.log`, or `ZERTO_HOOK_LOG`. One line per decision.
## Why the timeouts are derived, not fixed
## Known issue: false denials on cloud-protected VPGs
This hook used to deny every change to a cloud-protected VM.
The `win2019-1` denial above was a **false negative**, and it is worth
understanding before relying on this hook in an estate with cloud-protected
workloads.
`wait_for_tag` gave up after a hardcoded 45s. That is generous on a
`wait_for_tag` gives up after a hardcoded 45s. That is generous for a
vSphere-protected VPG, which checkpoints every 5s and surfaces a tag in about
4s, and impossible on an AWS-protected one, where a tag takes ~128s because the
journal only checkpoints every 630s.
4s. It is far too short elsewhere: a tag takes ~34s to appear on an
Azure-protected VPG and ~128s on an AWS-protected one, because journal cadence
is set by the protected site (5s vSphere, 60s Azure, 630s AWS).
So the guard reported "no checkpoint, refusing the change" while Zerto was in the
middle of creating one. The checkpoint landed a minute later, in the journal,
after the agent had already been told there was no rewind point.
So the hook denied the change, and the checkpoint landed anyway. It is in the
journal as `cp 56`, timestamped a minute after the hook reported failure. The
guard told the agent there was no rewind point while Zerto was in the middle of
creating one.
That is a worse failure than the one this hook exists to prevent. It is silent,
it looks correct in the log, and it blocks legitimate work on every
cloud-protected VM in the estate.
That failure mode is worse than the one the hook guards against, because it is
silent and looks correct: legitimate work is refused on every cloud-protected VM
while the log reads like the guard is doing its job.
Both budgets are now derived from the VPG's measured cadence rather than
guessed, which is why the numbers above differ by a factor of fifteen between
platforms.
Until `wait_for_tag` becomes cadence-aware, scope this hook's matcher to
vSphere-protected workloads.
+1 -1
View File
@@ -8,7 +8,7 @@
{
"type": "command",
"command": "/home/you/zerto-ai-rewind/.venv/bin/python /home/you/zerto-ai-rewind/hooks/zerto_guard_hook.py",
"timeout": 360
"timeout": 180
}
]
}
+2 -6
View File
@@ -24,7 +24,7 @@ tagged checkpoint and waits for the Zerto task to reach Completed.
"matcher": "mcp__.*",
"hooks": [{"type": "command",
"command": "/path/to/.venv/bin/python /path/to/hooks/zerto_guard_hook.py",
"timeout": 360}]}]}}
"timeout": 180}]}]}}
"""
from __future__ import annotations
@@ -40,11 +40,7 @@ LOG = os.environ.get("ZERTO_HOOK_LOG", os.path.expanduser("~/.zerto-guard-hook.l
# Seconds the hook will wait for the checkpoint. Must stay under the hook
# timeout configured in settings.json, or the host cancels us and the tool
# call proceeds unguarded through the normal permission flow.
# Must exceed the largest tag wait the guard can take. That is now derived from
# the VPG's checkpoint cadence and capped at 300s (MAX_TAG_TIMEOUT_S), because a
# tag takes ~128s to surface on an AWS-protected VPG. Too small a budget here
# just moves the false denial from the guard into the hook.
GUARD_TIMEOUT_S = float(os.environ.get("ZERTO_HOOK_GUARD_TIMEOUT", "330"))
GUARD_TIMEOUT_S = float(os.environ.get("ZERTO_HOOK_GUARD_TIMEOUT", "150"))
UNKNOWN_DECISION = os.environ.get("ZERTO_HOOK_UNKNOWN", "prompt") # prompt | allow | deny
+11 -73
View File
@@ -4,7 +4,6 @@ from __future__ import annotations
import asyncio
import re
import statistics
from datetime import UTC, datetime
from typing import Any
@@ -24,15 +23,6 @@ _WS = re.compile(r"\s+")
# so the name stays readable in the Zerto UI checkpoint list.
TAG_MAX_LEN = 250
# Tag-wait budget. Checkpoint cadence is set by the VPG's protected site and
# measured 5s (vSphere), 60s (Azure) and 630s (AWS) in one estate, so the wait
# has to be derived rather than fixed.
DEFAULT_TAG_TIMEOUT_S = 45.0 # cadence unmeasurable; the old constant
MIN_TAG_TIMEOUT_S = 45.0
MAX_TAG_TIMEOUT_S = 300.0 # AWS needed 128s; this leaves real headroom
MIN_POLL_INTERVAL_S = 1.5
MAX_POLL_INTERVAL_S = 15.0
def _clean(value: Any, limit: int) -> str:
"""One field of a checkpoint name: single-line, no separator collisions."""
@@ -84,81 +74,29 @@ def checkpoint_id(row: dict[str, Any]) -> str:
return str(value) if value is not None else ""
def checkpoint_gaps(rows: list[dict[str, Any]], sample: int = 20) -> list[float]:
"""Seconds between consecutive checkpoints, newest `sample` of them."""
stamps: list[datetime] = []
for row in rows[-sample:]:
raw = str(pick(row, "TimeStamp", "Timestamp", "timestamp") or "")
try:
stamps.append(datetime.fromisoformat(raw))
except ValueError:
continue
return [
(stamps[i + 1] - stamps[i]).total_seconds()
for i in range(len(stamps) - 1)
if stamps[i + 1] >= stamps[i]
]
def cadence_seconds(rows: list[dict[str, Any]], sample: int = 20) -> float | None:
"""How often this VPG writes a checkpoint. None when it cannot be measured."""
gaps = checkpoint_gaps(rows, sample)
return statistics.median(gaps) if gaps else None
def tag_wait_budget(cadence: float | None) -> tuple[float, float]:
"""How long to wait for a tag, and how often to look, given the cadence.
Cadence is set by the VPG's PROTECTED site, and the spread is enormous:
measured 5s on vSphere, 60s on Azure, 630s on AWS. A single constant cannot
serve all three. The old fixed 45s was generous for vSphere and impossible
for AWS, where a tag took 128s to surface, so the guard reported "no
checkpoint" while Zerto was still creating one and the change was refused
for no reason.
Visibility does not scale linearly with cadence (the insert makes its own
off-cadence checkpoint), so this is 2x cadence plus headroom, clamped.
"""
if cadence is None or cadence <= 0:
return DEFAULT_TAG_TIMEOUT_S, MIN_POLL_INTERVAL_S
timeout = min(max(2 * cadence + 30, MIN_TAG_TIMEOUT_S), MAX_TAG_TIMEOUT_S)
interval = min(max(cadence / 10, MIN_POLL_INTERVAL_S), MAX_POLL_INTERVAL_S)
return timeout, interval
async def wait_for_tag(
client: ZertoClient,
vpg_identifier: str,
tag: str,
*,
timeout_s: float | None = None,
interval_s: float | None = None,
timeout_s: float = 45.0,
interval_s: float = 1.5,
) -> dict[str, Any]:
"""Wait until the tag is listed. Budget derived from the VPG's own cadence."""
rows = await client.list_checkpoints(vpg_identifier)
for row in rows:
if checkpoint_tag(row) == tag:
return row
cadence = cadence_seconds(rows)
budget, poll = tag_wait_budget(cadence)
timeout_s = budget if timeout_s is None else timeout_s
interval_s = poll if interval_s is None else interval_s
deadline = asyncio.get_event_loop().time() + timeout_s
last: list[dict[str, Any]] = []
while asyncio.get_event_loop().time() < deadline:
await asyncio.sleep(interval_s)
rows = await client.list_checkpoints(vpg_identifier)
for row in rows:
last = await client.list_checkpoints(vpg_identifier)
for row in last:
if checkpoint_tag(row) == tag:
return row
measured = f"{cadence:.0f}s" if cadence else "unknown"
await asyncio.sleep(interval_s)
raise ZertoError(
f"Tagged checkpoint {tag!r} did not appear on VPG {vpg_identifier} "
f"within {timeout_s:.0f}s (this VPG checkpoints about every {measured}). "
"Do not mutate. Check the Zerto task before assuming the insert failed: "
"a completed task with no visible checkpoint means the wait was short, "
"not that the insert was rejected."
f"within {timeout_s:.0f}s. Do not mutate. "
"On a cloud-protected VPG the tag routinely takes longer than this to appear "
"(measured ~34s on Azure, ~128s on AWS), so this timeout may simply be too "
"short rather than the insert having failed. Check the Zerto task before "
"assuming it did not land."
)
+6 -43
View File
@@ -2,8 +2,6 @@
from __future__ import annotations
import base64
import hashlib
import json
from pathlib import Path
from typing import Any
@@ -34,20 +32,6 @@ _catalog: MutatingCatalog | None = None
_settings: dict[str, Any] = {}
# Files above this are refused rather than streamed through a tool result.
# FLR is for a config file or a dropped directory; a disk image belongs on the
# whole-VM ladder in docs/recover-ladder.md.
MAX_RECOVER_BYTES = 1_048_576
def _encode_recovered(blob: bytes) -> dict[str, Any]:
"""Text where it is text, base64 otherwise, so the caller can just use it."""
try:
return {"encoding": "text", "content": blob.decode("utf-8")}
except UnicodeDecodeError:
return {"encoding": "base64", "content": base64.b64encode(blob).decode("ascii")}
def _dump(payload: Any) -> str:
return json.dumps(payload, indent=2, default=str)
@@ -463,20 +447,16 @@ async def zerto_recover_file(
checkpoint_identifier: str,
guest_path: str,
confirmed: bool = False,
dest_dir: str | None = None,
) -> str:
"""File-level recovery from a journal checkpoint. The VM stays up.
Returns the file CONTENT, not a path. The caller may be on another machine,
so a path on this host is of no use to them, and letting a caller choose
where bytes land is an arbitrary write once this server is shared.
Requires confirmed=true (human yes). Cannot run during clone/test/live/EJC.
10.9 FLR Operator role fails; use an Administrator account.
Locally replicated VPGs only. FLR is performed at the VPG's recovery site,
so a VPG replicating to a cloud ZCA must be recovered from that ZCA's API.
guest_path is the path on the guest: /home/x/f.conf or C:\\Users\\x\\f.txt.
Files larger than max_recover_bytes are refused; use the whole-VM ladder.
"""
if not confirmed:
return _dump(
@@ -492,10 +472,7 @@ async def zerto_recover_file(
gate = await _flr_site_gate(client, vpg_identifier)
if not gate.get("ok"):
return _dump(gate)
# Server-owned, from config only. This string is also handed to the
# appliance as initialDownloadPath, so a caller-supplied value would let a
# caller point the ZVM at a path of their choosing.
dest = Path(_settings.get("recovery_dir") or "./recovered")
dest = Path(dest_dir or _settings.get("recovery_dir") or "./recovered")
dest.mkdir(parents=True, exist_ok=True)
session_id: str | None = None
before = await _live_session_ids(client)
@@ -514,29 +491,15 @@ async def zerto_recover_file(
token = download_token_from(token_payload)
blob = await client.fetch_download(token)
name = Path(guest_path.replace("\\", "/")).name or "recovered.bin"
limit = int(_settings.get("max_recover_bytes") or MAX_RECOVER_BYTES)
if len(blob) > limit:
payload = {
"ok": False,
"too_large": True,
"bytes": len(blob),
"limit": limit,
"message": (
f"{name} is {len(blob)} bytes, over the {limit} byte limit for "
"file level recovery. Use a bounded whole-VM operation instead: "
"offsite clone or failover test."
),
}
else:
out_path = dest / name
out_path.write_bytes(blob)
payload = {
"ok": True,
"name": name,
"path": str(out_path.resolve()),
"bytes": len(blob),
"sha256": hashlib.sha256(blob).hexdigest(),
"session_id": session_id,
"flr_path": flr_path,
**_encode_recovered(blob),
"message": f"Recovered {len(blob)} bytes of {name} from the journal.",
"message": f"Wrote {len(blob)} bytes to {out_path}",
}
except ZertoError as exc:
payload = {"ok": False, "session_id": session_id, "message": str(exc)}
-92
View File
@@ -1,7 +1,5 @@
from datetime import UTC, datetime
import pytest
from zerto_rewind_mcp.checkpoints import (
TAG_MAX_LEN,
checkpoint_id,
@@ -52,93 +50,3 @@ def test_checkpoint_row_keys():
row2 = {"checkpointId": "cp-2", "tag": "t"}
assert checkpoint_id(row2) == "cp-2"
assert checkpoint_tag(row2) == "t"
def _rows(*offsets_seconds):
from datetime import timedelta
base = datetime(2026, 9, 22, 12, 0, 0, tzinfo=UTC)
return [
{"TimeStamp": (base + timedelta(seconds=o)).isoformat().replace("+00:00", "Z")}
for o in offsets_seconds
]
def test_cadence_measures_the_median_gap():
from zerto_rewind_mcp.checkpoints import cadence_seconds
assert cadence_seconds(_rows(0, 5, 10, 15, 20)) == 5.0
assert cadence_seconds(_rows(0, 60, 120, 180)) == 60.0
# one irregular gap must not drag the answer around
assert cadence_seconds(_rows(0, 5, 10, 400, 405, 410)) == 5.0
def test_cadence_is_none_when_unmeasurable():
from zerto_rewind_mcp.checkpoints import cadence_seconds
assert cadence_seconds([]) is None
assert cadence_seconds([{"TimeStamp": "not-a-date"}]) is None
assert cadence_seconds(_rows(0)) is None # one checkpoint gives no gap
def test_tag_wait_budget_covers_every_measured_platform():
"""The three cadences measured in one estate, and what each actually needed."""
from zerto_rewind_mcp.checkpoints import tag_wait_budget
for cadence, observed_visibility in ((5.0, 4.0), (60.0, 34.0), (630.0, 128.0)):
budget, interval = tag_wait_budget(cadence)
assert budget > observed_visibility, (
f"cadence {cadence}s budgets {budget}s but the tag took {observed_visibility}s"
)
assert interval >= 1.5
def test_tag_wait_budget_is_clamped_at_both_ends():
from zerto_rewind_mcp.checkpoints import (
DEFAULT_TAG_TIMEOUT_S,
MAX_TAG_TIMEOUT_S,
MIN_TAG_TIMEOUT_S,
tag_wait_budget,
)
assert tag_wait_budget(0.1)[0] == MIN_TAG_TIMEOUT_S # absurdly fast VPG
assert tag_wait_budget(100_000)[0] == MAX_TAG_TIMEOUT_S # absurdly slow one
assert tag_wait_budget(None)[0] == DEFAULT_TAG_TIMEOUT_S
assert tag_wait_budget(None)[1] == 1.5
def test_wait_for_tag_returns_immediately_when_already_present():
import asyncio
from zerto_rewind_mcp.checkpoints import wait_for_tag
class Client:
def __init__(self):
self.calls = 0
async def list_checkpoints(self, vpg):
self.calls += 1
return [{"Tag": "ai:x", "CheckpointId": "7"}]
c = Client()
row = asyncio.run(wait_for_tag(c, "vpg", "ai:x"))
assert row["CheckpointId"] == "7"
assert c.calls == 1 # no sleep, no second poll
def test_wait_for_tag_error_names_the_measured_cadence():
import asyncio
from zerto_rewind_mcp.checkpoints import wait_for_tag
from zerto_rewind_mcp.client import ZertoError
class Client:
async def list_checkpoints(self, vpg):
return _rows(0, 60, 120, 180) # 60s cadence, tag never appears
with pytest.raises(ZertoError) as err:
# explicit tiny timeout so the test does not actually wait 150s
asyncio.run(wait_for_tag(Client(), "vpg", "ai:missing", timeout_s=0.01, interval_s=0.01))
msg = str(err.value)
assert "about every 60s" in msg
assert "was short, not that the insert was rejected" in msg
-41
View File
@@ -256,44 +256,3 @@ def test_flr_gate_refuses_remote_recovery_site_and_names_it():
# must tell the operator where the operation actually lives
assert out["recovery_site"] == "aws-zca"
assert "aws-zca" in out["message"]
def test_encode_recovered_text_and_binary():
from zerto_rewind_mcp.server import _encode_recovered
text = _encode_recovered(b"listen: 0.0.0.0:8443\n")
assert text["encoding"] == "text"
assert text["content"] == "listen: 0.0.0.0:8443\n"
binary = _encode_recovered(b"\x89PNG\r\n\x1a\n\xff\xfe")
assert binary["encoding"] == "base64"
import base64 as b64
assert b64.b64decode(binary["content"]) == b"\x89PNG\r\n\x1a\n\xff\xfe"
def test_recover_file_takes_no_caller_destination():
"""The caller must not choose where bytes land, nor where the ZVM mounts.
dest_dir used to be a tool parameter whose value was also passed to the
appliance as initialDownloadPath.
"""
import inspect
from zerto_rewind_mcp.server import zerto_recover_file
params = set(inspect.signature(zerto_recover_file).parameters)
assert "dest_dir" not in params
assert params == {
"vpg_identifier",
"vm_identifier",
"checkpoint_identifier",
"guest_path",
"confirmed",
}
def test_recover_byte_cap_is_configurable_and_has_a_default():
from zerto_rewind_mcp.server import MAX_RECOVER_BYTES
assert MAX_RECOVER_BYTES == 1_048_576