From 5b142e31925bd60380f86ca12e70e4e95ae03de7 Mon Sep 17 00:00:00 2001 From: igalshilman Date: Tue, 29 Sep 2026 19:09:28 +0200 Subject: [PATCH 1/2] Show a real kill-mid-turn recording in the README The README opens with a 54-second recording of the reference agent on the real stack: a turn calls three tools and starts a 30-second durable timer, the service is killed with kill -9 and started again, and Restate replays the journal so the turn finishes without re-running a recorded step. .misc/video records it unattended: record.mjs drives the service, the web UI and the Admin API journal onto director.html through a small CDP client with no npm dependencies, and encode.sh makes the MP4. The README there explains how to record a new take and upload it. Co-Authored-By: Claude Opus 5.5 (1M context) --- .misc/video/.gitignore | 2 + .misc/video/README.md | 70 +++++++ .misc/video/cdp.mjs | 167 ++++++++++++++++ .misc/video/director.html | 311 ++++++++++++++++++++++++++++++ .misc/video/encode.sh | 23 +++ .misc/video/record.mjs | 389 ++++++++++++++++++++++++++++++++++++++ README.md | 5 + 7 files changed, 967 insertions(+) create mode 100644 .misc/video/.gitignore create mode 100644 .misc/video/README.md create mode 100644 .misc/video/cdp.mjs create mode 100644 .misc/video/director.html create mode 100755 .misc/video/encode.sh create mode 100644 .misc/video/record.mjs diff --git a/.misc/video/.gitignore b/.misc/video/.gitignore new file mode 100644 index 0000000..3acb573 --- /dev/null +++ b/.misc/video/.gitignore @@ -0,0 +1,2 @@ +out/ +.DS_Store diff --git a/.misc/video/README.md b/.misc/video/README.md new file mode 100644 index 0000000..126733b --- /dev/null +++ b/.misc/video/README.md @@ -0,0 +1,70 @@ +# README video: kill it mid-turn + +A real recording of the reference agent: a turn calls three tools, starts a +30-second durable timer, the service is killed with `kill -9`, started again, +and Restate replays the turn's journal to finish it. The script records the +whole thing unattended, so it can be recorded again when the UI or the +runtime changes. + +| File | Does | +| --- | --- | +| `record.mjs` | Runs the scenario: starts, kills and restarts the core service, types into the web UI, polls the turn's journal from the Admin API, and captures frames of the stage | +| `director.html` | The 1920×1080 stage: the web UI, the service terminal, the journal panel, captions and title cards | +| `cdp.mjs` | A tiny Chrome DevTools Protocol client (Node 22+, no npm packages) | +| `encode.sh` | Turns the captured frames into `out/kill-it-mid-turn.mp4` | + +Everything on screen comes from the running system except the titles and +captions: the UI is streamed from its own tab, the terminal shows the +service's real log lines, and the journal panel is `sys_journal`. The quiet +stretches (the service down, the timer running) play at 4×, shown by the +badge in the header. + +## Record + +You need Docker (or `restate-server`), Google Chrome, ffmpeg, and an +`OPENAI_API_KEY`. Start from a fresh Restate server, so the agent ID `demo` +is new: + +```sh +pnpm build +docker run -d --name video-restate --rm -p 8080:8080 -p 9070:9070 \ + -e RESTATE_EXPERIMENTAL_ENABLE_PROTOCOL_V7=true docker.restate.dev/restatedev/restate:latest + +# Register the service once, then stop it: record.mjs starts it on camera. +(cd packages/libs/core && node dist/app.js) & +curl localhost:9070/deployments --json '{"uri":"http://host.docker.internal:9080"}' +kill %1 + +# The web UI. +(cd packages/apps/web && npx next start --hostname 127.0.0.1) & + +AGENT_ID=demo node .misc/video/record.mjs # about 90 seconds +.misc/video/encode.sh # out/kill-it-mid-turn.mp4 +``` + +With `restate-server` instead of Docker, register `http://localhost:9080`. + +## Check it + +Look at frames before using a take, for example one every three seconds: + +```sh +ffmpeg -i .misc/video/out/kill-it-mid-turn.mp4 -vf fps=1/3 /tmp/frames/%02d.png +``` + +Check that the ask went through at once (the prompt must not sit in the +composer), that the kill happened while only the timer was open, that the +journal marks only finished steps as replayed, and that the Restate UI scene +shows the 30-second `sleep` row. + +A take runs about 55 seconds. If it runs much longer, look for leftover +headless Chrome processes (`pgrep -f readme-video-chrome`): they slow the UI +down enough to delay the ask. + +## Put it in the README + +GitHub plays an MP4 inline only from a `user-attachments` URL: drag +`kill-it-mid-turn.mp4` into the README editor on github.com (or into a PR +comment) and use the URL it inserts. A committed `.mp4` file only renders as +a link. The README uses it under the intro; replace that URL when you record +a new take. diff --git a/.misc/video/cdp.mjs b/.misc/video/cdp.mjs new file mode 100644 index 0000000..0784824 --- /dev/null +++ b/.misc/video/cdp.mjs @@ -0,0 +1,167 @@ +// A minimal Chrome DevTools Protocol client: launches headless Chrome, opens +// tabs and sends commands to them. Node 22+ only (built-in WebSocket), so the +// recorder needs no npm packages. +import {spawn} from "node:child_process"; +import {mkdtempSync} from "node:fs"; +import {tmpdir} from "node:os"; +import {join} from "node:path"; + +const CHROME = + process.env.CHROME ?? "/Applications/Google Chrome.app/Contents/MacOS/Google Chrome"; + +export async function launchChrome({port = 9333} = {}) { + const profile = mkdtempSync(join(tmpdir(), "readme-video-chrome-")); + const chrome = spawn( + CHROME, + [ + "--headless=new", + `--remote-debugging-port=${port}`, + `--user-data-dir=${profile}`, + "--hide-scrollbars", + "--force-color-profile=srgb", + "--no-proxy-server", + "--no-first-run", + "--no-default-browser-check", + "about:blank", + ], + {stdio: "ignore"}, + ); + const url = await waitForDebugger(port); + const browser = new Browser(new WebSocket(url), chrome); + await browser.opened; + return browser; +} + +async function waitForDebugger(port) { + for (let attempt = 0; attempt < 100; attempt++) { + try { + const response = await fetch(`http://127.0.0.1:${port}/json/version`); + return (await response.json()).webSocketDebuggerUrl; + } catch { + await new Promise((resolve) => setTimeout(resolve, 100)); + } + } + throw new Error("Chrome did not start its debugger"); +} + +class Browser { + constructor(socket, process) { + this.socket = socket; + this.process = process; + this.nextId = 1; + this.waiting = new Map(); + this.listeners = new Map(); + this.opened = new Promise((resolve) => socket.addEventListener("open", resolve)); + socket.addEventListener("message", (message) => this.receive(JSON.parse(message.data))); + } + + receive(message) { + if (message.id !== undefined) { + const waiter = this.waiting.get(message.id); + this.waiting.delete(message.id); + if (message.error) waiter.reject(new Error(message.error.message)); + else waiter.resolve(message.result); + return; + } + const key = `${message.sessionId ?? ""}:${message.method}`; + for (const listener of this.listeners.get(key) ?? []) listener(message.params); + } + + send(method, params = {}, sessionId) { + const id = this.nextId++; + this.socket.send(JSON.stringify({id, method, params, sessionId})); + return new Promise((resolve, reject) => this.waiting.set(id, {resolve, reject})); + } + + on(sessionId, method, listener) { + const key = `${sessionId}:${method}`; + this.listeners.set(key, [...(this.listeners.get(key) ?? []), listener]); + } + + async newTab({url, width, height, scale = 1}) { + const {targetId} = await this.send("Target.createTarget", {url: "about:blank", newWindow: true}); + const {sessionId} = await this.send("Target.attachToTarget", {targetId, flatten: true}); + const tab = new Tab(this, sessionId); + await tab.send("Page.enable"); + await tab.send("Runtime.enable"); + await tab.send("Emulation.setDeviceMetricsOverride", { + width, + height, + deviceScaleFactor: scale, + mobile: false, + }); + if (url) await tab.goto(url); + return tab; + } + + /** Closes Chrome and waits for it to exit, so no renderers are left behind. */ + async close() { + const exited = new Promise((resolve) => this.process.once("exit", resolve)); + await this.send("Browser.close").catch(() => {}); + this.socket.close(); + await Promise.race([exited, new Promise((resolve) => setTimeout(resolve, 3_000))]); + this.process.kill("SIGKILL"); + } +} + +class Tab { + constructor(browser, sessionId) { + this.browser = browser; + this.sessionId = sessionId; + } + + send(method, params) { + return this.browser.send(method, params, this.sessionId); + } + + on(method, listener) { + this.browser.on(this.sessionId, method, listener); + } + + async goto(url) { + const loaded = new Promise((resolve) => this.on("Page.loadEventFired", resolve)); + await this.send("Page.navigate", {url}); + await loaded; + } + + async evaluate(expression) { + const {result, exceptionDetails} = await this.send("Runtime.evaluate", { + expression, + awaitPromise: true, + returnByValue: true, + }); + if (exceptionDetails) throw new Error(exceptionDetails.exception?.description ?? expression); + return result.value; + } + + /** Polls a page expression until it is truthy. */ + async waitFor(expression, {timeoutMs = 120_000, intervalMs = 200} = {}) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + const value = await this.evaluate(expression); + if (value) return value; + await new Promise((resolve) => setTimeout(resolve, intervalMs)); + } + throw new Error(`Timed out waiting for: ${expression}`); + } + + async screenshot() { + const {data} = await this.send("Page.captureScreenshot", {format: "png"}); + return Buffer.from(data, "base64"); + } + + /** Types text into the focused element, one key event per character. */ + async type(text, {delayMs = 35} = {}) { + for (const char of text) { + await this.send("Input.insertText", {text: char}); + await new Promise((resolve) => setTimeout(resolve, delayMs)); + } + } + + async press(key) { + const codes = {Enter: 13}; + const event = {key, code: key, windowsVirtualKeyCode: codes[key], text: key === "Enter" ? "\r" : undefined}; + await this.send("Input.dispatchKeyEvent", {type: "keyDown", ...event}); + await this.send("Input.dispatchKeyEvent", {type: "keyUp", ...event}); + } +} diff --git a/.misc/video/director.html b/.misc/video/director.html new file mode 100644 index 0000000..561f49c --- /dev/null +++ b/.misc/video/director.html @@ -0,0 +1,311 @@ + + + + + +Kill it mid-turn + + + +
+ +

Restate reference agent

+ — kill it mid-turn + ⏩ 4× + REAL RECORDING +
+ +
+
127.0.0.1:3000 · reference UI
+
+
+ +
+
agent service · packages/libs/coreRUNNING
+
+
+ +
+
turn journal · stored in Restate
+
+
+
+ +
+
localhost:9070 · Restate UI
+ +
+ +
1
+ +

+
+
+ + + + diff --git a/.misc/video/encode.sh b/.misc/video/encode.sh new file mode 100755 index 0000000..4b05cad --- /dev/null +++ b/.misc/video/encode.sh @@ -0,0 +1,23 @@ +#!/usr/bin/env bash +# Encodes the frames record.mjs wrote into an MP4 for GitHub. +# Usage: .misc/video/encode.sh [out dir] (default: .misc/video/out) +# Frames arrive only when the stage repaints, so each one is held until the +# next; ffmpeg then resamples that to a constant 30 fps. +set -euo pipefail +out="${1:-$(dirname "$0")/out}" + +node -e ' + const frames = require(process.argv[1] + "/frames.json"); + const lines = ["ffconcat version 1.0"]; + for (let i = 0; i < frames.length - 1; i++) { + lines.push(`file frames/${frames[i].file}`, `duration ${(frames[i + 1].time - frames[i].time).toFixed(4)}`); + } + lines.push(`file frames/${frames.at(-1).file}`); + require("fs").writeFileSync(process.argv[1] + "/frames.ffconcat", lines.join("\n") + "\n"); +' "$out" + +ffmpeg -y -loglevel error -f concat -safe 0 -i "$out/frames.ffconcat" \ + -vf "fps=30,format=yuv420p" -c:v libx264 -preset slow -crf 20 \ + -movflags +faststart "$out/kill-it-mid-turn.mp4" + +ls -lh "$out/kill-it-mid-turn.mp4" diff --git a/.misc/video/record.mjs b/.misc/video/record.mjs new file mode 100644 index 0000000..cd3aff4 --- /dev/null +++ b/.misc/video/record.mjs @@ -0,0 +1,389 @@ +// Records the README video: a real turn on the real stack, with the agent +// service killed while the turn waits on a durable timer. See README.md here. +// +// It starts and kills the core service itself, drives the web UI in headless +// Chrome, polls the turn's journal from the Restate Admin API, and composes +// everything on director.html. The director's screencast frames are written to +// /frames with the output time of each frame, for encode.sh. +import {spawn} from "node:child_process"; +import {mkdirSync, rmSync, writeFileSync} from "node:fs"; +import {dirname, join, resolve} from "node:path"; +import {createInterface} from "node:readline"; +import {fileURLToPath} from "node:url"; + +import {launchChrome} from "./cdp.mjs"; + +const HERE = dirname(fileURLToPath(import.meta.url)); +const REPO = resolve(HERE, "../.."); +const OUT = resolve(process.argv[2] ?? join(HERE, "out")); + +const AGENT_ID = process.env.AGENT_ID ?? `demo-${Date.now().toString(36)}`; +const UI_URL = process.env.UI_URL ?? "http://127.0.0.1:3000"; +const ADMIN_URL = process.env.RESTATE_ADMIN_URL ?? "http://localhost:9070"; +const SLEEP_SECONDS = 30; +const PROMPT = + `Get the weather in Berlin, Tokyo and New York in parallel. ` + + `Then sleep for ${SLEEP_SECONDS} seconds. Then compare the three cities in two sentences.`; + +const pause = (ms) => new Promise((done) => setTimeout(done, ms)); + +// --------------------------------------------------------------------------- +// Output timeline. Frames carry the wall-clock time they were painted; the +// speed segments map that to video time, so waiting stretches play faster. + +const speedSegments = [{at: Date.now() / 1000, factor: 1}]; +const frames = []; + +function videoTime(seconds) { + let time = 0; + for (let i = 0; i < speedSegments.length; i++) { + const segment = speedSegments[i]; + const end = speedSegments[i + 1]?.at ?? Infinity; + if (seconds <= segment.at) break; + time += (Math.min(seconds, end) - segment.at) / segment.factor; + } + return time; +} + +async function setSpeed(stage, factor) { + speedSegments.push({at: Date.now() / 1000, factor}); + await stage.evaluate(`stage.speed(${factor})`); +} + +// --------------------------------------------------------------------------- +// The agent service: started, killed and restarted by this script, its log +// lines shown in the terminal panel. + +let service; +let onServiceLine = () => {}; + +function startService() { + service = spawn("node", ["./dist/app.js"], { + cwd: join(REPO, "packages/libs/core"), + env: {...process.env, NODE_ENV: "production"}, + stdio: ["ignore", "pipe", "pipe"], + }); + for (const stream of [service.stdout, service.stderr]) { + createInterface({input: stream}).on("line", (line) => onServiceLine(line)); + } + return service.pid; +} + +// "[restate][2026-...Z][AgentSession/demo/doTurn][inv_...] INFO: Replaying invocation." +function formatServiceLine(line) { + const match = line.match(/^\[restate\]\[[^\]]*T([\d:]+)\.\d+Z\](?:\[([^\]]+)\]\[[^\]]+\])? \w+: (.*)$/); + if (!match) return null; + const [, time, target, message] = match; + if (!target) return message.startsWith("Restate SDK started") ? {text: `${time} listening on :9080`} : null; + const handler = target.split("/").pop(); + if (!["doTurn", "ask"].includes(handler)) return null; + const service = target.startsWith("AgentSession") ? "AgentSession" : "Agent"; + const text = `${time} ${service}.${handler} ${message}`; + return {text, kind: message.startsWith("Replaying") ? "hot" : ""}; +} + +// --------------------------------------------------------------------------- +// The turn's journal, read from the Restate Admin API. + +async function query(sql) { + const response = await fetch(`${ADMIN_URL}/query`, { + method: "POST", + headers: {"content-type": "application/json", accept: "application/json"}, + body: JSON.stringify({query: sql}), + }); + return (await response.json()).rows ?? []; +} + +async function currentTurnId() { + const rows = await query( + `SELECT id FROM sys_invocation WHERE target_service_name = 'AgentSession' ` + + `AND target_handler_name = 'doTurn' AND target_service_key = '${AGENT_ID}' ` + + `ORDER BY created_at DESC LIMIT 1`, + ); + return rows[0]?.id; +} + +async function journalEntries(turnId) { + return query( + `SELECT index, entry_type, name FROM sys_journal WHERE id = '${turnId}' ` + + `AND entry_type IN ('Command: Input', 'Command: Run', 'Command: Sleep', 'Command: Output', ` + + `'Notification: Run', 'Notification: Sleep') ORDER BY index`, + ); +} + +const LABELS = { + "agent-model": "model call", + getWeather: "getWeather", + "discover-agent-tools": "tool catalog", +}; + +// Turns journal entries into panel rows. A run counts as recorded once as many +// run results as runs up to it are in the journal; parallel runs finish within +// milliseconds of each other, so the order does not show. +function journalRows(entries, replay) { + const rows = []; + const runResults = entries.filter((entry) => entry.entry_type === "Notification: Run").length; + const sleepDone = entries.some((entry) => entry.entry_type === "Notification: Sleep"); + let runs = 0; + for (const entry of entries) { + const row = {index: entry.index, label: "", name: "", state: "wait", status: "running…"}; + if (entry.entry_type === "Command: Input") { + Object.assign(row, {label: "input", name: "the message", state: "done", status: "✓ recorded"}); + } else if (entry.entry_type === "Command: Output") { + Object.assign(row, {label: "output", name: "turn finished", state: "done", status: "✓ recorded"}); + } else if (entry.entry_type === "Command: Run") { + runs++; + const recorded = runs <= runResults; + const name = LABELS[entry.name] ?? entry.name; + const isTool = !LABELS[entry.name] || entry.name === "getWeather"; + Object.assign(row, { + label: isTool ? "tool" : "", + name, + state: recorded ? "done" : "wait", + status: recorded ? "✓ recorded" : "running…", + }); + } else if (entry.entry_type === "Command: Sleep") { + Object.assign(row, { + label: "timer", + name: `sleep ${SLEEP_SECONDS}s`, + state: sleepDone ? "done" : "wait", + status: sleepDone ? "✓ fired" : "⏱ in Restate", + }); + } else { + continue; + } + if (replay.replayed.has(entry.index)) Object.assign(row, {state: "replayed", status: "♻ replayed"}); + else if (replay.crashedAt !== undefined && entry.index > replay.crashedAt && row.state === "done") { + Object.assign(row, {state: "fresh", status: "✓ new"}); + } + rows.push(row); + } + return rows; +} + +// --------------------------------------------------------------------------- + +async function main() { + rmSync(OUT, {recursive: true, force: true}); + mkdirSync(join(OUT, "frames"), {recursive: true}); + + const browser = await launchChrome(); + const stage = await browser.newTab({url: `file://${HERE}/director.html`, width: 1920, height: 1080}); + // Both UIs render zoomed in, so they stay legible in a README-sized player. + // The web UI's agent picker is cropped off the top: the conversation is the + // part the video is about. + const uiScale = 1100 / 900; + // The UI opens only once the service runs (its default agent is "demo" too). + // Opened earlier, its calls to the agent would wait in Restate's retry + // backoff, and the first ask behind them. + const ui = await browser.newTab({url: "about:blank", width: 900, height: 900, scale: uiScale}); + async function openUi() { + await ui.goto(`${UI_URL}/?agent=${AGENT_ID}`); + const pickerHeight = await ui.waitFor( + `[...document.querySelectorAll("button")].find((b) => b.textContent.trim() === "Open agent")?.getBoundingClientRect().bottom`, + ); + const cropCss = pickerHeight + 13; + await ui.send("Emulation.setDeviceMetricsOverride", { + width: 900, + height: Math.round(768 / uiScale + cropCss), + deviceScaleFactor: uiScale, + mobile: false, + }); + await stage.evaluate(`stage.uiCrop(${Math.round(cropCss * uiScale)})`); + } + const restateScale = 1.25; + const restateUi = await browser.newTab({ + url: "about:blank", + width: Math.round(1840 / restateScale), + height: Math.round(768 / restateScale), + scale: restateScale, + }); + + // Stream a tab into the director. Only the newest frame is forwarded. + function mirror(tab, method) { + let latest; + let busy = false; + tab.on("Page.screencastFrame", async ({data, sessionId}) => { + tab.send("Page.screencastFrameAck", {sessionId}); + latest = data; + if (busy) return; + busy = true; + while (latest) { + const frame = latest; + latest = undefined; + await stage.evaluate(`stage.${method}(${JSON.stringify(frame)})`); + } + busy = false; + }); + return tab.send("Page.startScreencast", {format: "jpeg", quality: 92, maxWidth: 4096, maxHeight: 4096}); + } + + let frameNumber = 0; + stage.on("Page.screencastFrame", ({data, metadata, sessionId}) => { + stage.send("Page.screencastFrameAck", {sessionId}); + const file = `${String(frameNumber++).padStart(6, "0")}.jpg`; + writeFileSync(join(OUT, "frames", file), Buffer.from(data, "base64")); + frames.push({file, time: videoTime(metadata.timestamp)}); + }); + + await mirror(ui, "uiFrame"); + await mirror(restateUi, "restateFrame"); + + const say = (expression) => stage.evaluate(expression); + const js = JSON.stringify; + + onServiceLine = (line) => { + const formatted = formatServiceLine(line); + if (formatted) say(`stage.line(${js(formatted.text)}, ${js(formatted.kind ?? "")})`); + }; + + async function typeCommand(command) { + await say("stage.prompt()"); + for (const char of command) { + await say(`stage.typeChar(${js(char)})`); + await pause(38); + } + await pause(350); + } + + const replay = {replayed: new Set(), crashedAt: undefined}; + let turnId; + let entries = []; + let polling = true; + const poller = (async () => { + while (polling) { + turnId ??= await currentTurnId(); + if (turnId) { + entries = await journalEntries(turnId); + await say(`stage.journal(${js(journalRows(entries, replay))})`); + } + await pause(250); + } + })(); + + // --- Title --------------------------------------------------------------- + await say(`stage.card("Restate reference agent", "Kill it mid-turn.", "A real turn, a real crash, no lost work.")`); + await stage.send("Page.startScreencast", {format: "jpeg", quality: 90}); + await pause(3200); + await say("stage.hideCard()"); + + // --- 1. Start the service and ask --------------------------------------- + await say(`stage.caption(1, "Start the agent and ask for something slow", "Three tools in parallel, then a ${SLEEP_SECONDS}-second durable timer.")`); + await typeCommand("node dist/app.js"); + let pid = startService(); + await say(`stage.service("running", "RUNNING · PID ${pid}")`); + await pause(800); + await openUi(); + await pause(700); + + // The UI shows "Live" once its notification stream to the agent is up. + await ui.waitFor(`document.body.innerText.includes("Live")`); + const composer = `document.querySelector('textarea[aria-label="Ask message"]')`; + await ui.evaluate(`${composer}.focus()`); + await ui.type(PROMPT, {delayMs: 22}); + await pause(400); + await ui.press("Enter"); + await ui.waitFor(`${composer}.value === ""`, {timeoutMs: 30_000}); + + // --- 2. Journal fills ----------------------------------------------------- + while (!turnId) await pause(100); + await say(`stage.caption(2, "Every step lands in the turn's journal", "Model calls and tool results are stored in Restate as they finish.")`); + const quiet = () => { + const sleeping = entries.some((entry) => entry.entry_type === "Command: Sleep"); + const runs = entries.filter((entry) => entry.entry_type === "Command: Run").length; + const results = entries.filter((entry) => entry.entry_type === "Notification: Run").length; + return sleeping && runs === results; + }; + while (!quiet()) await pause(200); + await pause(2500); + // The model may take one more step while it waits; kill only when no run is open. + while (!quiet()) await pause(200); + + // --- 3. Kill ---------------------------------------------------------------- + await say(`stage.caption(3, "Kill the service mid-turn", "kill -9: no shutdown, no warning. The timer is still running.")`); + await typeCommand(`kill -9 ${pid}`); + replay.crashedAt = entries.at(-1).index; + service.kill("SIGKILL"); + await say("stage.crash()"); + await say(`stage.service("down", "KILLED")`); + await say(`stage.line("[1]+ Killed: 9 node dist/app.js", "bad")`); + await pause(2200); + await say(`stage.caption(3, "No process is running this turn", "Restate holds its journal and its timer. Nothing is lost.")`); + await setSpeed(stage, 4); + await pause(8000); + await setSpeed(stage, 1); + + // --- 4. Restart and replay ------------------------------------------------ + await say(`stage.caption(4, "Start it again", "Restate replays the journal. Recorded results are reused, not re-run.")`); + await typeCommand("node dist/app.js"); + await say(`stage.service("starting", "STARTING")`); + const replaying = new Promise((done) => { + const previous = onServiceLine; + onServiceLine = (line) => { + previous(line); + if (line.includes("doTurn") && line.includes("Replaying invocation")) done(); + }; + }); + pid = startService(); + await replaying; + await say(`stage.service("running", "RUNNING · PID ${pid}")`); + + // Only finished steps are replayed from their recorded result; the timer is + // still pending in Restate and keeps its own status. + const recorded = journalRows(entries, replay).filter( + (row) => row.index <= replay.crashedAt && row.state === "done", + ); + for (const row of recorded) { + replay.replayed.add(row.index); + await say(`stage.journal(${js(journalRows(entries, replay))})`); + await pause(140); + } + const models = recorded.filter((row) => row.name === "model call").length; + const tools = recorded.filter((row) => row.label === "tool").length; + await say( + `stage.journalFoot(${js(`Replayed ${models} model calls and ${tools} tool calls. None ran again.`)})`, + ); + + // --- 5. The turn finishes --------------------------------------------------- + await pause(2500); + await say(`stage.caption(5, "The turn finishes where it stopped", "The timer fires on schedule and the model writes its answer.")`); + await setSpeed(stage, 4); + while (!entries.some((entry) => entry.entry_type === "Command: Output")) await pause(200); + await setSpeed(stage, 1); + await pause(4000); + + // --- 6. The Restate UI ---------------------------------------------------- + await restateUi.goto(`${ADMIN_URL}/ui/invocations/${turnId}`); + await pause(2500); + await say(`stage.caption(6, "Inspect it in the Restate UI", "One invocation, every step of the turn, across the crash.")`); + await say("stage.showRestate(true)"); + await pause(2500); + // Scroll to the timer: its bar spans the time the service was down. + await restateUi.evaluate(`(() => { + const timer = [...document.querySelectorAll("span")].find((span) => span.textContent.trim() === "sleep"); + const top = timer.getBoundingClientRect().top + scrollY - innerHeight / 2; + scrollTo({top, behavior: "smooth"}); + })()`); + await pause(2000); + await say(`stage.caption(6, "The timer kept running while the service was down", "It fired on schedule, and the turn went on from the next step.")`); + await pause(4000); + + // --- End card --------------------------------------------------------------- + await say(`stage.card("github.com/restatedev/agent", "Durable by default.", "Every model call, tool and timer — journaled by Restate.")`); + await pause(3500); + + polling = false; + await poller; + await stage.send("Page.stopScreencast"); + frames.push({file: frames.at(-1).file, time: videoTime(Date.now() / 1000)}); + writeFileSync(join(OUT, "frames.json"), JSON.stringify(frames)); + service.kill(); + await browser.close(); + console.log(`${frames.length} frames, ${frames.at(-1).time.toFixed(1)}s, agent ${AGENT_ID}, turn ${turnId}`); +} + +main().catch((error) => { + console.error(error); + service?.kill("SIGKILL"); + process.exit(1); +}); diff --git a/README.md b/README.md index ef314cc..4bde79a 100644 --- a/README.md +++ b/README.md @@ -4,6 +4,11 @@ A complete agent, built on [Restate](https://restate.dev). Every feature a modern agent needs is here as a small module you can read in one sitting, and Restate keeps each turn running through crashes and days-long waits. +https://github.com/user-attachments/assets/5929008d-736e-4cda-9466-d5bb2b5a8245 + +*A real recording: the service is killed mid-turn, started again, and the +turn finishes without repeating a model call or a tool.* + [Features](#features) · [One turn, start to finish](#one-turn-start-to-finish) · [How a turn works](#how-a-turn-works) · [Quickstart](#quickstart) · [Documentation](docs/README.md) From 32fd7209eb428b578d68d5cb8870c940a46ab4e3 Mon Sep 17 00:00:00 2001 From: igalshilman Date: Tue, 29 Sep 2026 19:37:39 +0200 Subject: [PATCH 2/2] Record one session: an ordinary turn first, then the crash The README video now opens with an ordinary turn: the model writes a program whose three web searches run at once, saves a plan in its sandbox, and is steered while it works. Only then does a follow-up wait on a durable timer, get killed with kill -9 and finish after the replay. The recorder drives both acts in one agent: the journal panel follows the newest turn, the transcript keeps following new events, and the take stops if the steer reaches no running turn. Co-Authored-By: Claude Opus 5.5 (1M context) --- .misc/video/README.md | 42 +++++--- .misc/video/director.html | 8 +- .misc/video/encode.sh | 4 +- .misc/video/record.mjs | 217 +++++++++++++++++++++++++------------- README.md | 7 +- 5 files changed, 178 insertions(+), 100 deletions(-) diff --git a/.misc/video/README.md b/.misc/video/README.md index 126733b..5a28b6d 100644 --- a/.misc/video/README.md +++ b/.misc/video/README.md @@ -1,23 +1,29 @@ -# README video: kill it mid-turn +# README video: a real session -A real recording of the reference agent: a turn calls three tools, starts a -30-second durable timer, the service is killed with `kill -9`, started again, -and Restate replays the turn's journal to finish it. The script records the -whole thing unattended, so it can be recorded again when the UI or the -runtime changes. +A real recording of the reference agent, in two acts: + +1. **An ordinary turn.** A casual request for a weekend plan. The model + writes a program whose three web searches run at once, saves `lisbon.md` + in its sandbox, and is steered while it works. +2. **Kill it mid-turn.** A follow-up waits on a 30-second durable timer. The + service is killed with `kill -9`, started again, and Restate replays the + turn's journal, so the turn reads the file back and answers. + +The script records the whole thing unattended, so it can be recorded again +when the UI or the runtime changes. | File | Does | | --- | --- | | `record.mjs` | Runs the scenario: starts, kills and restarts the core service, types into the web UI, polls the turn's journal from the Admin API, and captures frames of the stage | | `director.html` | The 1920×1080 stage: the web UI, the service terminal, the journal panel, captions and title cards | | `cdp.mjs` | A tiny Chrome DevTools Protocol client (Node 22+, no npm packages) | -| `encode.sh` | Turns the captured frames into `out/kill-it-mid-turn.mp4` | +| `encode.sh` | Turns the captured frames into `out/agent-demo.mp4` | Everything on screen comes from the running system except the titles and captions: the UI is streamed from its own tab, the terminal shows the service's real log lines, and the journal panel is `sys_journal`. The quiet -stretches (the service down, the timer running) play at 4×, shown by the -badge in the header. +stretches (model calls after the steer, the service down, the timer +running) play faster, shown by the badge in the header. ## Record @@ -38,8 +44,8 @@ kill %1 # The web UI. (cd packages/apps/web && npx next start --hostname 127.0.0.1) & -AGENT_ID=demo node .misc/video/record.mjs # about 90 seconds -.misc/video/encode.sh # out/kill-it-mid-turn.mp4 +AGENT_ID=demo node .misc/video/record.mjs # about two minutes +.misc/video/encode.sh # out/agent-demo.mp4 ``` With `restate-server` instead of Docker, register `http://localhost:9080`. @@ -49,22 +55,24 @@ With `restate-server` instead of Docker, register `http://localhost:9080`. Look at frames before using a take, for example one every three seconds: ```sh -ffmpeg -i .misc/video/out/kill-it-mid-turn.mp4 -vf fps=1/3 /tmp/frames/%02d.png +ffmpeg -i .misc/video/out/agent-demo.mp4 -vf fps=1/3 /tmp/frames/%02d.png ``` -Check that the ask went through at once (the prompt must not sit in the -composer), that the kill happened while only the timer was open, that the +Check that each ask went through at once (the prompt must not sit in the +composer), that the steer landed while the turn ran (the script stops if it +did not), that the kill happened while only the timer was open, that the journal marks only finished steps as replayed, and that the Restate UI scene -shows the 30-second `sleep` row. +shows the 30-second `sleep` row. The model's wording differs between takes, +so read its answers too. -A take runs about 55 seconds. If it runs much longer, look for leftover +A take runs about 75 seconds. If it runs much longer, look for leftover headless Chrome processes (`pgrep -f readme-video-chrome`): they slow the UI down enough to delay the ask. ## Put it in the README GitHub plays an MP4 inline only from a `user-attachments` URL: drag -`kill-it-mid-turn.mp4` into the README editor on github.com (or into a PR +`agent-demo.mp4` into the README editor on github.com (or into a PR comment) and use the URL it inserts. A committed `.mp4` file only renders as a link. The README uses it under the intro; replace that URL when you record a new take. diff --git a/.misc/video/director.html b/.misc/video/director.html index 561f49c..d9e271e 100644 --- a/.misc/video/director.html +++ b/.misc/video/director.html @@ -8,7 +8,7 @@ -Kill it mid-turn +Agent demo