diff --git a/.gitignore b/.gitignore index 321c656..edb21bb 100644 --- a/.gitignore +++ b/.gitignore @@ -16,3 +16,13 @@ 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 diff --git a/demo/README.md b/demo/README.md new file mode 100644 index 0000000..0e88d7a --- /dev/null +++ b/demo/README.md @@ -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= +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. diff --git a/demo/demo_win.example.json b/demo/demo_win.example.json new file mode 100644 index 0000000..fc68406 --- /dev/null +++ b/demo/demo_win.example.json @@ -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" +} diff --git a/demo/diagram.py b/demo/diagram.py new file mode 100644 index 0000000..55c9396 --- /dev/null +++ b/demo/diagram.py @@ -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) diff --git a/demo/driver.py b/demo/driver.py new file mode 100644 index 0000000..4df1bf2 --- /dev/null +++ b/demo/driver.py @@ -0,0 +1,241 @@ +"""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 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} $ cat {rec['path']}{RST}") + body = Path(rec["path"]).read_text() + 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']}" + 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 []) + \ + [rec["path"], f"{target}:{CFG['guest_path']}"] + w(f"{YEL} $ scp {Path(rec['path']).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())) diff --git a/demo/journal.py b/demo/journal.py new file mode 100644 index 0000000..f09cbfc --- /dev/null +++ b/demo/journal.py @@ -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()) diff --git a/demo/mux_vo.py b/demo/mux_vo.py new file mode 100644 index 0000000..7537074 --- /dev/null +++ b/demo/mux_vo.py @@ -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}") diff --git a/demo/narration.md b/demo/narration.md new file mode 100644 index 0000000..178b445 --- /dev/null +++ b/demo/narration.md @@ -0,0 +1,53 @@ +# 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, and +there's no other copy of that file anywhere. + +[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, and the only other copy is in last night's backup, which is already hours +behind. + +[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. diff --git a/demo/record.sh b/demo/record.sh new file mode 100755 index 0000000..ca8bcf3 --- /dev/null +++ b/demo/record.sh @@ -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" diff --git a/demo/record_win.sh b/demo/record_win.sh new file mode 100755 index 0000000..ade61b7 --- /dev/null +++ b/demo/record_win.sh @@ -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" diff --git a/demo/vo/build_narration.py b/demo/vo/build_narration.py new file mode 100644 index 0000000..71fcd2e --- /dev/null +++ b/demo/vo/build_narration.py @@ -0,0 +1,33 @@ +"""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") + +REPLACE = {"VM": "vee em", "ad1": "A D one"} # customer-facing script: few acronyms left +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") diff --git a/demo/windows_driver.py b/demo/windows_driver.py new file mode 100644 index 0000000..2233b69 --- /dev/null +++ b/demo/windows_driver.py @@ -0,0 +1,193 @@ +"""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 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 = Path(rec["path"]).read_text() + 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 = Path(rec["path"]).read_bytes() + 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()))