From 2a35397e99a9b23605471576a5274685ce4e7716 Mon Sep 17 00:00:00 2001 From: Justin Paul Date: Mon, 21 Sep 2026 21:17:24 -0400 Subject: [PATCH] demo: commit the recording harness The demo tooling only existed in a session scratch directory, which is temporary. This puts it in the repo so the video can be rebuilt. Terminal only: tmux drives a two pane session, asciinema records it, agg renders it, ffmpeg encodes it. The left pane runs the loop against a live ZVM, the right pane polls the VPG journal so the tagged checkpoint appears on camera as it lands. Narration is synthesised per beat and aligned to marks the driver writes while it runs, rather than to predicted timings, so a slow API call or an FLR mount that takes longer than usual does not drift the audio. Two things that has to respect are written down in the README: agg's idle-time-limit must exceed the longest pause or it compresses idle time and breaks the wall-clock mapping, and each beat holds for its narration length so no line is cut off. demo_win.json carries live guest credentials, so only an example with placeholders is committed and the real file is gitignored, along with the generated wav, cast, gif and mp4. The xAI key path and voice id come from the environment now instead of being hardcoded to one machine. Also records the narration gotchas that cost time: mapping an acronym to run-together phonetics ({"VM": "vee em"}) is spoken as one word, "vem"; prose written for the page sounds robotic read aloud; and volumedetect reports no samples when aimed at a file whose first stream is video, which makes a working audio track look silent. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_016yVfC5nvZowoLFnEGWhLGn --- .gitignore | 10 ++ demo/README.md | 75 ++++++++++++ demo/demo_win.example.json | 12 ++ demo/diagram.py | 44 +++++++ demo/driver.py | 241 +++++++++++++++++++++++++++++++++++++ demo/journal.py | 78 ++++++++++++ demo/mux_vo.py | 39 ++++++ demo/narration.md | 53 ++++++++ demo/record.sh | 52 ++++++++ demo/record_win.sh | 42 +++++++ demo/vo/build_narration.py | 33 +++++ demo/windows_driver.py | 193 +++++++++++++++++++++++++++++ 12 files changed, 872 insertions(+) create mode 100644 demo/README.md create mode 100644 demo/demo_win.example.json create mode 100644 demo/diagram.py create mode 100644 demo/driver.py create mode 100644 demo/journal.py create mode 100644 demo/mux_vo.py create mode 100644 demo/narration.md create mode 100755 demo/record.sh create mode 100755 demo/record_win.sh create mode 100644 demo/vo/build_narration.py create mode 100644 demo/windows_driver.py 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()))