mirror of
https://github.com/obra/superpowers.git
synced 2026-09-28 21:25:03 +08:00
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:
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user