fix(movie): own recorder cleanup and bound observation

Implement Task 3 of the consolidated PR 2214 movie repairs. Startup failures could leak an already launched ttyd or browser because acquisitions preceded the cleanup block; close killed historical numeric PIDs and reported success without confirmed cleanup. Register each acquired resource inside serve's try/finally, invalidate readiness on shutdown, retain logs, and atomically retire PID metadata only after confirmed owned-process and profile cleanup. Close now requests stop and waits under one overall 30-second deadline without killing PIDs.

Poll CDP and owner readiness while observing commands without recording, preserving a completed native exit status across simultaneous disconnects and reserving exit 2 for live unfinished commands. Fill only bounded 5 fps slots when a final capture crosses the hard or hold endpoint, correctly strip ST-terminated OSC titles, and emit ASCII-escaped stdout JSON while preserving UTF-8 session files. A serve-only CDP call-boundary check prevents a stop from waiting through multiple consecutive calls.

Add 21 contract tests using fake processes, fake websocket responses, fake clocks, byte-token capture callbacks, and intercepted frame writes. Register real-session test cleanup immediately after Popen and retain live PID snapshots, but do not execute that fixture. RED evidence and implementation report are in .superpowers/sdd/2026-09-11-movie-committee-repairs/task-3-report.md. Validation: 26/26 authorized prompt, serve-argument, and recorder contract tests pass; git diff --check passes. No media generation, inspection, browser/ttyd launch, real session or frame-grid suite was performed. Native movie acceptance remains Drew's review; an abruptly hard-killed owner with a live browser and stale readiness remains the agreed limitation.
This commit is contained in:
Drew Ritter
2026-09-11 20:01:08 -07:00
parent e48cb57bf1
commit 42486dbb6d
3 changed files with 617 additions and 56 deletions
@@ -115,7 +115,7 @@ def prompt_command(kind, cwd):
f"[Convert]::FromBase64String('{encoded}'))))")
VISIBLE = re.compile(rb"\x1b\][^\x07\x1b]*(?:\x07|\x1b\\\\)|\x1b\[[0-?]*[ -/]*[@-~]|\r")
VISIBLE = re.compile(rb"\x1b\][^\x07\x1b]*(?:\x07|\x1b\\)|\x1b\[[0-?]*[ -/]*[@-~]|\r")
def at_prompt(log):
@@ -153,7 +153,13 @@ def film(out, seconds, hold, capture, finished, clock=time.monotonic, sleep=time
now = clock()
if stop is None and finished():
stop = now + hold
if now - start >= seconds or (stop is not None and now >= stop):
endpoint = min(start + seconds, stop) if stop is not None else start + seconds
if now >= endpoint:
# A capture may cross the endpoint. Fill only grid slots before
# that endpoint; tolerate floating-point noise at exact FPS ticks.
while last is not None and (index + 1) / FPS < endpoint - start - 1e-9:
index += 1
(out / f"f{index:05d}.png").write_bytes(last)
return index + 1
slot = int((now - start) * FPS)
if slot > index:
@@ -170,7 +176,7 @@ class CDP:
import websocket
self.ws = websocket.create_connection(url, timeout=5, suppress_origin=True)
self.count, self.on_event = 0, None
self.count, self.on_event, self.before_call = 0, None, None
def recv(self, timeout):
import websocket
@@ -188,6 +194,8 @@ class CDP:
return message
def call(self, method, params=None, timeout=10):
if self.before_call:
self.before_call()
self.count += 1
self.ws.send(json.dumps({"id": self.count, "method": method, "params": params or {}}))
deadline = time.monotonic() + timeout
@@ -259,7 +267,7 @@ def lit_fraction(png):
def serve(args):
directory = args.session
directory = args.session.resolve()
if directory.exists() and any(directory.iterdir()):
raise SystemExit(f"{directory} is not empty: use a new session directory")
directory.mkdir(parents=True, exist_ok=True)
@@ -286,19 +294,21 @@ def serve(args):
"--disable-background-networking", "--remote-debugging-address=127.0.0.1",
f"--remote-debugging-port={debug_port}", f"--user-data-dir={directory / 'profile'}",
f"--window-size={WIDTH},{HEIGHT}", "--hide-scrollbars", "about:blank"]
logs = [(directory / "ttyd.log").open("ab"), (directory / "browser.log").open("ab")]
processes = [
subprocess.Popen(ttyd_argv, cwd=cwd, stdin=subprocess.DEVNULL, stdout=logs[0],
stderr=subprocess.STDOUT, start_new_session=unix),
subprocess.Popen(browser_argv, cwd=directory, stdin=subprocess.DEVNULL, stdout=logs[1],
stderr=subprocess.STDOUT, start_new_session=unix),
]
logs, processes = [], []
session = dict(shell=args.shell, cwd=str(cwd), terminal_url=f"http://127.0.0.1:{port}/",
debug_port=debug_port, pids=[process.pid for process in processes])
write_json(directory / "session.json", session)
output = (directory / "terminal.log").open("ab")
debug_port=debug_port, pids=[])
output, cdp = None, None
state = {"closed": False}
class StopRequested(Exception):
pass
def check_active():
if (directory / "stop").exists():
raise StopRequested()
if state["closed"] or any(process.poll() is not None for process in processes):
raise ConnectionError("the terminal session closed")
def on_event(event):
if event["method"] == "Network.webSocketFrameReceived":
frame = event["params"]["response"]
@@ -312,9 +322,9 @@ def serve(args):
def pump_until(condition, timeout, failure):
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
check_active()
cdp.recv(0.05)
if state["closed"]:
raise ConnectionError("the terminal connection closed")
check_active()
if condition():
return
(directory / "timeout.png").write_bytes(screenshot(cdp))
@@ -325,16 +335,28 @@ def serve(args):
# it has been silent for `quiet` seconds.
deadline, seen = time.monotonic() + quiet, output.tell()
while time.monotonic() < deadline:
check_active()
cdp.recv(0.05)
if state["closed"]:
raise ConnectionError("the terminal connection closed")
check_active()
if output.tell() != seen:
deadline, seen = time.monotonic() + quiet, output.tell()
code = 1
try:
for name in ("ttyd.log", "browser.log"):
logs.append((directory / name).open("ab"))
processes.append(subprocess.Popen(
ttyd_argv, cwd=cwd, stdin=subprocess.DEVNULL, stdout=logs[0],
stderr=subprocess.STDOUT, start_new_session=unix))
processes.append(subprocess.Popen(
browser_argv, cwd=directory, stdin=subprocess.DEVNULL, stdout=logs[1],
stderr=subprocess.STDOUT, start_new_session=unix))
session["pids"] = [process.pid for process in processes]
write_json(directory / "session.json", session)
output = (directory / "terminal.log").open("ab")
deadline = time.monotonic() + 20
while True:
check_active()
try:
cdp = CDP(page_url(debug_port))
break
@@ -343,6 +365,7 @@ def serve(args):
raise TimeoutError(f"browser did not start: {error}") from None
time.sleep(0.1)
cdp.on_event = on_event
cdp.before_call = check_active
cdp.call("Network.enable") # before navigation, or the terminal socket is never reported
cdp.call("Page.enable")
cdp.call("Emulation.setDeviceMetricsOverride",
@@ -385,29 +408,71 @@ def serve(args):
raise RuntimeError(f"the terminal renders blank ({lit:.4%} lit pixels); see ready.png")
typed("clear")
time.sleep(0.3)
check_active()
prompt = prompts(tail(directory / "terminal.log"))[-1]
write_json(directory / "ready.json", dict(session, prompt=prompt, lit=round(lit, 4)))
print(json.dumps({"ready": True, "session": str(directory), "cwd": prompt["cwd"]},
ensure_ascii=False), flush=True)
print(json.dumps({"ready": True, "session": str(directory), "cwd": prompt["cwd"]}), flush=True)
while not (directory / "stop").exists():
cdp.recv(0.2)
if state["closed"]:
raise ConnectionError("the terminal connection closed")
check_active()
code = 0
except KeyboardInterrupt:
except (KeyboardInterrupt, StopRequested):
code = 0
except Exception as error: # noqa: BLE001 - report, then clean up below
print(f"serve: {error}", file=sys.stderr)
code = 0 if (directory / "stop").exists() else 1
code = 1
finally:
for process in processes:
kill_process_tree(process.pid)
cleaned = True
try:
(directory / "ready.json").unlink(missing_ok=True)
except OSError as error:
print(f"serve cleanup: {error}", file=sys.stderr)
cleaned = False
if cdp is not None:
try:
cdp.ws.close()
except Exception as error:
print(f"serve cleanup: {error}", file=sys.stderr)
cleaned = False
for process in processes:
try:
# Only the owner acts on handles it acquired, never stored PIDs.
if process.poll() is None:
kill_process_tree(process.pid)
process.wait(timeout=5)
except subprocess.TimeoutExpired:
pass
for handle in (output, *logs):
handle.close()
except (OSError, subprocess.TimeoutExpired) as error:
print(f"serve cleanup: {error}", file=sys.stderr)
cleaned = False
for handle in ([output] if output is not None else []) + logs:
try:
handle.close()
except OSError as error:
print(f"serve cleanup: {error}", file=sys.stderr)
cleaned = False
deadline = time.monotonic() + 5
while True:
try:
shutil.rmtree(directory / "profile")
break
except FileNotFoundError:
break
except OSError as error:
if time.monotonic() >= deadline:
print(f"serve cleanup: {error}", file=sys.stderr)
cleaned = False
break
time.sleep(0.1)
if cleaned:
try:
# Readers see either active ownership or completed cleanup.
completed = directory / "session.json.tmp"
write_json(completed, dict(session, pids=[], closed=True))
completed.replace(directory / "session.json")
except OSError as error:
print(f"serve cleanup: {error}", file=sys.stderr)
cleaned = False
if not cleaned:
code = 1
return code
@@ -418,14 +483,37 @@ def observe(args, cdp, n0):
def latest():
return next((p for p in reversed(prompts(tail(log))) if p["n"] > n0), None)
def poll():
prompt = latest()
if prompt:
return prompt
try:
cdp.recv(0.05)
except Exception:
# The owner can log the final native status just before disconnect.
prompt = latest()
if prompt:
return prompt
raise
prompt = latest()
if prompt:
return prompt
if not (args.session / "ready.json").is_file():
raise ConnectionError("the terminal session is no longer ready")
return None
frames = 0
if args.record:
frames = film(args.record, args.seconds, args.hold, lambda: screenshot(cdp), lambda: latest() is not None)
else:
deadline = time.monotonic() + args.seconds
while latest() is None and time.monotonic() < deadline:
time.sleep(0.05)
prompt = latest()
try:
if args.record:
frames = film(args.record, args.seconds, args.hold, lambda: screenshot(cdp), lambda: poll() is not None)
else:
deadline = time.monotonic() + args.seconds
while poll() is None and time.monotonic() < deadline:
time.sleep(0.05)
prompt = poll()
except Exception as error:
print(json.dumps({"outcome": "failed", "error": str(error)}))
return 1
result = {"outcome": "completed" if prompt else "running"}
if prompt:
result.update(ok=prompt["ok"], exit_code=prompt["exit_code"], cwd=prompt["cwd"])
@@ -433,7 +521,7 @@ def observe(args, cdp, n0):
result["frames"] = frames
result["scene"] = {"kind": "frames", "src": str(args.record.resolve()), "rate": FPS}
write_json(args.record / "take.json", result)
print(json.dumps(result, ensure_ascii=False))
print(json.dumps(result))
return 2 if not prompt else 0 if prompt["ok"] else 1
@@ -471,20 +559,24 @@ def watch(args):
def close(args):
session = read_json(args.session / "session.json")
(args.session / "stop").write_text("")
for pid in session["pids"]:
kill_process_tree(pid)
for _ in range(50): # the browser releases its profile shortly after dying
try:
shutil.rmtree(args.session / "profile")
break
except FileNotFoundError:
break
except OSError:
time.sleep(0.1)
print(json.dumps({"closed": True}))
return 0
deadline = time.monotonic() + 30
try:
session = read_json(args.session / "session.json")
(args.session / "stop").write_text("", encoding="utf-8")
while True:
if (session.get("closed") is True and session.get("pids") == []
and not (args.session / "ready.json").exists()
and not (args.session / "profile").exists()):
print(json.dumps({"closed": True}))
return 0
remaining = deadline - time.monotonic()
if remaining <= 0:
raise TimeoutError("serve did not confirm cleanup within 30 seconds")
time.sleep(min(0.1, remaining))
session = read_json(args.session / "session.json")
except (OSError, ValueError, AttributeError) as error:
print(f"close: {error}", file=sys.stderr)
return 1
def main():
@@ -0,0 +1,462 @@
"""Recorder decisions with fake processes/CDP and byte-token captures only."""
import contextlib
import importlib.util
import io
import json
import subprocess
import tempfile
import unittest
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import patch
SCRIPT = Path(__file__).resolve().parents[2] / "skills/proving-it-works-with-a-movie/examples/film-terminal.py"
def recorder():
spec = importlib.util.spec_from_file_location("recorder_contract", SCRIPT)
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
class Clock:
def __init__(self):
self.now = 0.0
def monotonic(self):
return self.now
def sleep(self, seconds):
self.now += seconds
class Process:
def __init__(self, pid):
self.pid = pid
self.returncode = None
def poll(self):
return self.returncode
def wait(self, timeout):
if self.returncode is None:
raise subprocess.TimeoutExpired("fake child", timeout)
return self.returncode
@contextlib.contextmanager
def serving(failure=None, stop_at=None, relative=False):
module, clock = recorder(), Clock()
with tempfile.TemporaryDirectory() as temp, contextlib.ExitStack() as stack:
root = Path(temp)
directory = root / "session"
args = SimpleNamespace(session=directory, cwd=root, shell="bash", shell_exe=None,
ttyd="fake-ttyd", browser="fake-browser")
if relative:
import os
args.session = Path(os.path.relpath(directory))
handles, children, launches, tokens = [], [], [], []
original_open, original_write = Path.open, module.write_json
stopped = False
def open_file(path, *positional, **kwargs):
if positional == ("ab",):
if failure == path.name:
raise OSError("injected " + path.name)
handle = original_open(path, *positional, **kwargs)
handles.append(handle)
return handle
return original_open(path, *positional, **kwargs)
def popen(argv, **kwargs):
index = len(children)
if failure == ("ttyd launch" if index == 0 else "browser launch"):
raise OSError("injected launch")
child = Process(1100 + index)
children.append(child)
launches.append((argv, kwargs))
if index == 1:
(directory / "profile").mkdir()
return child
def kill(pid):
for child in children:
if child.pid == pid and failure != "child wait":
child.returncode = -9
def write_json(path, value):
if failure == "session metadata" and path.name == "session.json":
raise OSError("injected metadata failure")
original_write(path, value)
def stop(phase):
nonlocal stopped
if stop_at == phase and not stopped:
(directory / "stop").write_text("")
stopped = True
class FakeCDP:
def __init__(self, url):
if failure == "connection":
raise RuntimeError("injected connection failure")
self.n = 0
self.on_event = None
self.ws = SimpleNamespace(close=lambda: None)
def recv(self, timeout):
clock.sleep(timeout)
if (directory / "ready.json").exists():
stop("ready")
if failure == "later connection":
raise ConnectionError("injected later disconnect")
else:
stop("pump" if self.n == 0 else "quiet")
return None
def call(self, method, params=None, **kwargs):
if getattr(self, "before_call", None):
self.before_call()
if stop_at == "call":
stop("call")
clock.sleep(10)
if method == "Input.dispatchKeyEvent" and params["type"] == "keyDown":
self.n += 1
return {}
cdp = None
def connect(url):
nonlocal cdp
cdp = FakeCDP(url)
return cdp
def page_url(port):
stop("retry")
if stop_at == "retry":
raise OSError("not listening yet")
return "fake-url"
def tail(path):
if cdp and cdp.n:
return f"\x1b]0;MOVIE;{cdp.n};1;0;{root}\x07$ ".encode()
return b"$ "
for obj, name, value in (
(module, "shell_argv", lambda *a: ["fake-bash"]),
(module, "find_browser", lambda *a: "fake-browser"),
(module, "free_port", lambda: 1234),
(module.subprocess, "Popen", popen),
(module, "kill_process_tree", kill),
(module, "page_url", page_url),
(module, "CDP", connect),
(module, "write_json", write_json),
(module, "tail", tail),
(module, "screenshot", lambda *a: b"capture-token"),
(module, "lit_fraction", lambda *a: 0.01),
(module.time, "monotonic", clock.monotonic),
(module.time, "sleep", clock.sleep),
(Path, "open", open_file),
(Path, "write_bytes", lambda path, token: tokens.append((path.name, token))),
):
stack.enter_context(patch.object(obj, name, value))
if failure == "profile removal":
stack.enter_context(patch.object(module.shutil, "rmtree", side_effect=PermissionError("locked")))
stack.enter_context(contextlib.redirect_stdout(io.StringIO()))
stack.enter_context(contextlib.redirect_stderr(io.StringIO()))
try:
yield SimpleNamespace(module=module, args=args, directory=directory, children=children,
handles=handles, launches=launches, clock=clock)
finally:
# Release the test's log handles even when ownership assertions fail.
for handle in handles:
handle.close()
class ServeLifecycleTests(unittest.TestCase):
def test_every_acquisition_failure_releases_owned_resources(self):
for failure in ("ttyd.log", "browser.log", "ttyd launch", "browser launch",
"session metadata", "terminal.log", "connection", "later connection"):
with self.subTest(failure=failure), serving(failure) as rig:
try:
code = rig.module.serve(rig.args)
except Exception as error:
code = error
self.assertEqual(code, 1)
self.assertTrue(all(p.poll() is not None for p in rig.children))
self.assertTrue(all(h.closed for h in rig.handles))
self.assertFalse((rig.directory / "profile").exists())
self.assertFalse((rig.directory / "ready.json").exists())
for name in ("ttyd.log", "browser.log", "terminal.log"):
if any(Path(h.name).name == name for h in rig.handles):
self.assertTrue((rig.directory / name).exists(), "logs are evidence")
def test_relative_session_path_matches_browser_profile_and_cwd(self):
with serving("connection", relative=True) as rig:
self.assertEqual(rig.module.serve(rig.args), 1)
argv, kwargs = rig.launches[1]
profile_arg = next(a.split("=", 1)[1] for a in argv if a.startswith("--user-data-dir="))
self.assertEqual(Path(kwargs["cwd"]) / profile_arg, (rig.directory / "profile").resolve())
self.assertTrue(Path(kwargs["cwd"]).is_absolute())
def test_stop_requests_interrupt_startup_and_finish_owned_cleanup(self):
for phase in ("retry", "pump", "quiet", "ready"):
with self.subTest(phase=phase), serving(stop_at=phase) as rig:
self.assertEqual(rig.module.serve(rig.args), 0)
self.assertLess(rig.clock.now, 5)
self.assertTrue(all(p.poll() is not None for p in rig.children))
self.assertTrue(all(h.closed for h in rig.handles))
session = rig.module.read_json(rig.directory / "session.json")
self.assertEqual(session["pids"], [])
self.assertTrue(session["closed"])
self.assertFalse((rig.directory / "ready.json").exists())
self.assertFalse((rig.directory / "profile").exists())
def test_stop_during_one_startup_call_prevents_another_bounded_call(self):
with serving(stop_at="call") as rig:
self.assertEqual(rig.module.serve(rig.args), 0)
self.assertLessEqual(rig.clock.now, 10)
self.assertTrue(rig.module.read_json(rig.directory / "session.json")["closed"])
def test_pid_retirement_is_atomic_and_follows_resource_cleanup(self):
with serving(stop_at="ready") as rig:
replace = Path.replace
observed = []
def publish(path, target):
session = rig.module.read_json(target)
observed.append(session["pids"])
self.assertEqual(session["pids"], [1100, 1101])
self.assertTrue(all(p.poll() is not None for p in rig.children))
self.assertTrue(all(h.closed for h in rig.handles))
self.assertFalse((rig.directory / "profile").exists())
self.assertFalse((rig.directory / "ready.json").exists())
return replace(path, target)
with patch.object(Path, "replace", publish):
self.assertEqual(rig.module.serve(rig.args), 0)
self.assertEqual(observed, [[1100, 1101]])
self.assertEqual(rig.module.read_json(rig.directory / "session.json")["pids"], [])
def test_failed_cleanup_does_not_retire_owned_ids_or_claim_closed(self):
for failure in ("child wait", "profile removal"):
with self.subTest(failure=failure), serving(failure, stop_at="ready") as rig:
self.assertEqual(rig.module.serve(rig.args), 1)
session = rig.module.read_json(rig.directory / "session.json")
self.assertEqual(session["pids"], [1100, 1101])
self.assertFalse(session.get("closed", False))
self.assertFalse((rig.directory / "ready.json").exists())
class CloseContractTests(unittest.TestCase):
def test_close_waits_for_owner_cleanup_and_repeated_close_never_kills(self):
module, clock = recorder(), Clock()
with tempfile.TemporaryDirectory() as temp:
directory = Path(temp)
module.write_json(directory / "session.json", {"pids": [987654]})
(directory / "ready.json").write_text("{}")
(directory / "profile").mkdir()
(directory / "terminal.log").write_text("retain evidence")
def sleep(seconds):
clock.sleep(seconds)
self.assertTrue((directory / "stop").exists())
if clock.now >= 12:
(directory / "profile").rmdir()
(directory / "ready.json").unlink()
module.write_json(directory / "session.json", {"pids": [], "closed": True})
with patch.object(module, "kill_process_tree", side_effect=AssertionError("historical PID kill")), \
patch.object(module.time, "monotonic", clock.monotonic), \
patch.object(module.time, "sleep", sleep), contextlib.redirect_stdout(io.StringIO()):
self.assertEqual(module.close(SimpleNamespace(session=directory)), 0)
self.assertGreaterEqual(clock.now, 12)
self.assertEqual(module.close(SimpleNamespace(session=directory)), 0)
self.assertEqual((directory / "terminal.log").read_text(), "retain evidence")
def test_incomplete_or_unavailable_cleanup_is_not_success(self):
for state in ("missing", "unreadable", "unavailable owner", "profile remains", "ready remains"):
with self.subTest(state=state), tempfile.TemporaryDirectory() as temp:
module, clock, directory = recorder(), Clock(), Path(temp)
if state != "missing":
module.write_json(directory / "session.json", {"pids": [987654]})
if state == "unreadable":
(directory / "session.json").write_text("invalid json")
if state in ("profile remains", "ready remains"):
module.write_json(directory / "session.json", {"pids": [], "closed": True})
if state == "profile remains":
(directory / "profile").mkdir()
else:
(directory / "ready.json").write_text("{}")
with patch.object(module, "kill_process_tree", lambda pid: None), \
patch.object(module.time, "monotonic", clock.monotonic), \
patch.object(module.time, "sleep", clock.sleep), \
contextlib.redirect_stdout(io.StringIO()) as out, \
contextlib.redirect_stderr(io.StringIO()):
try:
code = module.close(SimpleNamespace(session=directory))
except Exception as error:
code = error
self.assertEqual(code, 1)
self.assertNotIn('"closed": true', out.getvalue())
self.assertLessEqual(clock.now, 30.1)
if state == "unavailable owner":
self.assertGreaterEqual(clock.now, 30)
class CDPCallTests(unittest.TestCase):
def test_stop_during_response_prevents_the_next_call_from_being_sent(self):
module, sent, stopped = recorder(), [], []
def recv():
stopped.append(True)
return json.dumps({"id": len(sent), "result": {}})
ws = SimpleNamespace(settimeout=lambda timeout: None, recv=recv,
send=lambda data: sent.append(json.loads(data)["method"]))
websocket = SimpleNamespace(create_connection=lambda *a, **kw: ws,
WebSocketTimeoutException=TimeoutError)
def check_active():
if stopped:
raise InterruptedError("stop requested")
with patch.dict("sys.modules", websocket=websocket):
cdp = module.CDP("fake-url")
cdp.before_call = check_active
self.assertEqual(cdp.call("Network.enable"), {})
with self.assertRaises(InterruptedError):
cdp.call("Page.enable")
self.assertEqual(sent, ["Network.enable"])
class ObservationTests(unittest.TestCase):
def observe(self, action, completed=False, seconds=0.2, cp1252=False):
module, clock = recorder(), Clock()
with tempfile.TemporaryDirectory() as temp:
directory = Path(temp)
log = directory / "terminal.log"
initial = b"\x1b]0;MOVIE;1;1;0;/work\x07"
final = "\x1b]0;MOVIE;2;0;7;C:/René/λ\x1b\\".encode()
log.write_bytes(initial + (final if completed else b""))
(directory / "ready.json").write_text("{}")
module.write_json(directory / "session.json", {"pids": [1100, 1101]})
calls = []
def recv(timeout):
calls.append(timeout)
clock.sleep(timeout)
if action == "disconnect":
raise ConnectionError("browser connection closed")
if action == "session loss":
(directory / "ready.json").unlink(missing_ok=True)
if action == "prompt arrives":
log.write_bytes(initial + final)
if action == "prompt then disconnect":
log.write_bytes(initial + final)
raise ConnectionError("browser connection closed")
return None
cdp = SimpleNamespace(recv=recv)
args = SimpleNamespace(session=directory, record=None, seconds=seconds, hold=0)
raw = io.BytesIO()
out = io.TextIOWrapper(raw, encoding="cp1252" if cp1252 else "utf-8")
with patch.object(module.time, "monotonic", clock.monotonic), \
patch.object(module.time, "sleep", clock.sleep), contextlib.redirect_stdout(out):
code = module.observe(args, cdp, 1)
out.flush()
result = json.loads(raw.getvalue().decode("ascii" if cp1252 else "utf-8"))
return code, result, calls
def test_disconnected_browser_is_failure_without_recording(self):
code, result, calls = self.observe("disconnect")
self.assertEqual(code, 1)
self.assertEqual(result["outcome"], "failed")
self.assertTrue(calls)
def test_lost_terminal_session_is_failure_even_with_live_browser(self):
code, result, calls = self.observe("session loss")
self.assertEqual(code, 1)
self.assertEqual(result["outcome"], "failed")
self.assertTrue(calls)
def test_only_live_unfinished_command_returns_two(self):
for seconds in (0, 0.2):
with self.subTest(seconds=seconds):
code, result, calls = self.observe("alive", seconds=seconds)
self.assertEqual((code, result), (2, {"outcome": "running"}))
self.assertTrue(calls, "even an expired observation must establish liveness")
def test_new_prompt_retains_native_failure_status(self):
for action in ("prompt arrives", "prompt then disconnect"):
with self.subTest(action=action):
code, result, _ = self.observe(action)
self.assertEqual(code, 1)
self.assertEqual(result, {"outcome": "completed", "ok": False,
"exit_code": 7, "cwd": "C:/René/λ"})
def test_completed_prompt_is_not_consumed_by_later_disconnect(self):
code, result, _ = self.observe("disconnect", completed=True)
self.assertEqual((code, result["outcome"], result["exit_code"]), (1, "completed", 7))
def test_stdout_json_round_trips_non_ascii_paths_on_cp1252(self):
try:
code, result, _ = self.observe("alive", completed=True, cp1252=True)
except UnicodeError as error:
self.fail(f"stdout JSON was not portable: {error}")
self.assertEqual((code, result["cwd"]), (1, "C:/René/λ"))
class VisibleTextTests(unittest.TestCase):
def test_visible_prompt_survives_bel_and_st_title_markers(self):
module = recorder()
for terminator in (b"\x07", b"\x1b\\"):
with self.subTest(terminator=terminator):
log = b"\x1b[32m/work $ \x1b]0;MOVIE;2;1;0;/work" + terminator
self.assertTrue(module.at_prompt(log))
self.assertFalse(module.at_prompt(log + b"busy"))
class CaptureTimingTests(unittest.TestCase):
def film(self, durations, seconds=1, hold=0, complete_at=None):
module, clock, writes, shots = recorder(), Clock(), [], []
def capture():
token = bytes([len(shots) + 1])
clock.sleep(durations[len(shots)] if len(shots) < len(durations) else 0)
shots.append(token)
return token
with tempfile.TemporaryDirectory() as temp, \
patch.object(Path, "write_bytes", lambda path, token: writes.append((path.name, token))):
frames = module.film(Path(temp), seconds, hold, capture,
lambda: complete_at is not None and clock.now >= complete_at,
clock.monotonic, clock.sleep)
self.assertEqual(len(writes), frames)
return frames, writes
def test_capture_crossing_hard_endpoint_fills_exactly_five_slots(self):
frames, writes = self.film([1.2])
self.assertEqual(frames, 5)
self.assertEqual(writes, [("f00000.png", b"\x01"), ("f00001.png", b"\x01"),
("f00002.png", b"\x01"), ("f00003.png", b"\x01"),
("f00004.png", b"\x01")])
def test_capture_crossing_hold_endpoint_is_bounded(self):
frames, writes = self.film([1.2], seconds=10, hold=0.6, complete_at=0)
self.assertEqual(frames, 3)
self.assertEqual([name for name, _ in writes], ["f00000.png", "f00001.png", "f00002.png"])
def test_mid_capture_stall_repeats_previous_token_without_missing_slots(self):
frames, writes = self.film([0.01, 0.5], seconds=1)
self.assertEqual(frames, 5)
self.assertEqual([token for _, token in writes], [b"\x01", b"\x02", b"\x02", b"\x03", b"\x04"])
def test_completion_and_hold_keep_the_normal_grid(self):
frames, _ = self.film([], seconds=10, hold=0.4, complete_at=1)
self.assertEqual(frames, 7)
def test_completion_without_hold_and_zero_duration_do_not_add_slots(self):
self.assertEqual(self.film([], seconds=10, complete_at=0)[0], 0)
self.assertEqual(self.film([], seconds=0)[0], 0)
if __name__ == "__main__":
unittest.main()
@@ -194,11 +194,15 @@ class SessionTests(unittest.TestCase):
"--cwd", str(self.work), "--ttyd", TTYD, "--browser", BROWSER]
if os.environ.get("MOVIE_TEST_SHELL_EXE"):
argv += ["--shell-exe", os.environ["MOVIE_TEST_SHELL_EXE"]]
self.owned_pids = []
self.serve = subprocess.Popen(argv, stdout=self.log, stderr=subprocess.STDOUT)
self.addCleanup(self.close_session)
deadline = time.monotonic() + 45
while not (self.session / "ready.json").exists() and self.serve.poll() is None \
and time.monotonic() < deadline:
time.sleep(0.1)
if (self.session / "session.json").exists():
self.owned_pids = json.loads((self.session / "session.json").read_text(encoding="utf-8"))["pids"]
if not (self.session / "ready.json").exists():
report = "".join(f"--- {name}\n" + path.read_text(errors="replace") if path.exists() else ""
for name, path in (("serve.log", Path(self.tmp.name) / "serve.log"),
@@ -206,15 +210,18 @@ class SessionTests(unittest.TestCase):
("browser.log", self.session / "browser.log")))
self.fail(report)
def tearDown(self):
def close_session(self):
if self.serve.poll() is None:
self.cli("close", str(self.session))
# A failed setup may not have session metadata yet, but the
# owner still needs its stop request before we wait for cleanup.
self.session.mkdir(parents=True, exist_ok=True)
(self.session / "stop").write_text("", encoding="utf-8")
try:
self.serve.wait(15)
self.serve.wait(35)
except subprocess.TimeoutExpired:
self.serve.kill()
self.serve.wait()
for pid in json.loads((self.session / "session.json").read_text(encoding="utf-8"))["pids"]:
for pid in self.owned_pids:
self.assertTrue(gone(pid), f"pid {pid} survived close")
def cli(self, *args, timeout=120):