demo: commit the recording harness #8
+10
@@ -16,3 +16,13 @@ catalog.json
|
|||||||
recovered/
|
recovered/
|
||||||
.claude/*
|
.claude/*
|
||||||
!.claude/gitea-ship.json
|
!.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
|
||||||
|
|||||||
@@ -0,0 +1,75 @@
|
|||||||
|
# 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.
|
||||||
@@ -0,0 +1,12 @@
|
|||||||
|
{
|
||||||
|
"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"
|
||||||
|
}
|
||||||
@@ -0,0 +1,44 @@
|
|||||||
|
"""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
@@ -0,0 +1,254 @@
|
|||||||
|
"""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()))
|
||||||
@@ -0,0 +1,78 @@
|
|||||||
|
"""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())
|
||||||
@@ -0,0 +1,39 @@
|
|||||||
|
"""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}")
|
||||||
@@ -0,0 +1,52 @@
|
|||||||
|
# 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.
|
||||||
Executable
+52
@@ -0,0 +1,52 @@
|
|||||||
|
#!/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"
|
||||||
Executable
+42
@@ -0,0 +1,42 @@
|
|||||||
|
#!/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"
|
||||||
@@ -0,0 +1,35 @@
|
|||||||
|
"""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")
|
||||||
@@ -0,0 +1,204 @@
|
|||||||
|
"""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()))
|
||||||
Reference in New Issue
Block a user