zerto_recover_file no longer writes the file to this host and hands back a path, because a path here means nothing to a caller elsewhere and letting the caller choose it was an arbitrary write. Both drivers still read rec["path"], so they broke. They now use the returned content. A small recovered_bytes() helper in each driver decodes the text or base64 form, so the Windows copy-back keeps shipping exact bytes rather than letting PowerShell rewrite line endings, which is the bug that put a stray CR in an earlier take. The Linux driver writes the bytes to a local file before scp, since scp needs something on disk to send. Verified against a live recovery: the helper returns 164 bytes whose sha256 matches the one the server reported, and the keys the on-screen show() filter uses are all still present. Co-Authored-By: Claude Opus 5 (1M context) <[email protected]> Claude-Session: https://claude.ai/code/session_016yVfC5nvZowoLFnEGWhLGn
205 lines
8.2 KiB
Python
205 lines
8.2 KiB
Python
"""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()))
|