From ed5bda02491c70518b4dac080b290f9cfc5ef020 Mon Sep 17 00:00:00 2001 From: Mike Arpaia Date: Thu, 24 Sep 2026 15:37:23 -0600 Subject: [PATCH 01/10] Add cooperative live-session shutdown and port reuse tests --- .github/workflows/live-shutdown-windows.yml | 70 ++++ docs/protocols/live-viewer-v1.md | 18 ++ python/src/microsimulator/viewer_server.py | 151 +++++++-- python/tests/test_viewer_shutdown.py | 336 ++++++++++++++++++++ viewer/README.md | 23 +- viewer/browser/live-shutdown.mjs | 174 ++++++++++ viewer/index.html | 8 + viewer/src/live.ts | 64 +++- viewer/src/main.ts | 19 +- viewer/src/style.css | 7 + viewer/tests/live.test.ts | 105 +++++- 11 files changed, 933 insertions(+), 42 deletions(-) create mode 100644 .github/workflows/live-shutdown-windows.yml create mode 100644 python/tests/test_viewer_shutdown.py create mode 100644 viewer/browser/live-shutdown.mjs diff --git a/.github/workflows/live-shutdown-windows.yml b/.github/workflows/live-shutdown-windows.yml new file mode 100644 index 0000000..d81fff7 --- /dev/null +++ b/.github/workflows/live-shutdown-windows.yml @@ -0,0 +1,70 @@ +name: Windows live-session shutdown + +on: + pull_request: + paths: + - .github/workflows/live-shutdown-windows.yml + - CMakeLists.txt + - cpp/** + - python/** + - pyproject.toml + - uv.lock + push: + branches: [master, marpaia/17] + workflow_dispatch: + +permissions: + contents: read + +concurrency: + group: windows-live-shutdown-${{ github.ref }} + cancel-in-progress: true + +jobs: + shutdown: + runs-on: windows-2025 + timeout-minutes: 20 + defaults: + run: + shell: pwsh + env: + CMAKE_BUILD_PARALLEL_LEVEL: "2" + # Exercise the host reference engine without requiring a GPU toolkit. + CMAKE_ARGS: -DCM_ENABLE_METAL=OFF -DCM_ENABLE_CUDA=OFF -DCM_BUILD_TESTS=OFF + steps: + - name: Check out source + uses: actions/checkout@v6 + with: + persist-credentials: false + + - name: Set up uv + uses: astral-sh/setup-uv@c771a70e6277c0a99b617c7a806ffedaca235ff9 # v9.0.0 + with: + enable-cache: true + python-version: "3.12" + + - name: Build CPU extension and install locked test dependencies + run: uv sync --locked --group dev + + - name: Record platform, shell, interpreter, and backend + run: | + New-Item -ItemType Directory -Force build/shutdown-evidence | Out-Null + $PSVersionTable | Out-File build/shutdown-evidence/platform.txt + [System.Environment]::OSVersion | Out-File -Append build/shutdown-evidence/platform.txt + uv run --no-sync python --version 2>&1 | Tee-Object -Append build/shutdown-evidence/platform.txt + uv run --no-sync microsimulator devices --json | Tee-Object -Append build/shutdown-evidence/platform.txt + + - name: Verify authenticated Stop, worker draining, and same-port process restart + run: >- + uv run --no-sync python -m pytest + python/tests/test_viewer_server.py python/tests/test_viewer_shutdown.py + -v --tb=short --junitxml=build/shutdown-evidence/tests.xml + + - name: Upload shutdown evidence + if: ${{ always() }} + uses: actions/upload-artifact@v7 + with: + name: windows-live-shutdown-${{ github.sha }} + path: build/shutdown-evidence + if-no-files-found: warn + retention-days: 30 diff --git a/docs/protocols/live-viewer-v1.md b/docs/protocols/live-viewer-v1.md index fbbf780..e2559bf 100644 --- a/docs/protocols/live-viewer-v1.md +++ b/docs/protocols/live-viewer-v1.md @@ -37,6 +37,15 @@ Rejected commands and model failures return a data-only error: { "type": "error", "message": "reason" } ``` +Intentional shutdown sends lifecycle notifications before closing the WebSocket with code 1000: + +```json +{"type":"session","state":"stopping"} +{"type":"session","state":"stopped"} +``` + +`stopping` means admission of new work has ended. `stopped` means the active operation has finished and the simulation worker has terminated. Clients should process previously received messages (including asynchronous scene verification) before interpreting the subsequent socket close. A close without `stopped` is still an unexpected disconnect. These messages extend the v1 vocabulary; use a viewer built from the same release as the server. + ## Client commands The vocabulary is closed. Unknown fields are rejected. @@ -48,8 +57,17 @@ The vocabulary is closed. Unknown fields are rejected. {"type":"pause"} {"type":"reset"} {"type":"checkpoint"} +{"type":"stop"} ``` `steps` defaults to one and is bounded to 1 through 10,000. Playback advances the configured number of steps per published frame. A step or reset first pauses playback. Disconnecting the final client pauses the simulation. Reset calls the original server-side model factory again with its original backend, device, seed, parameters, and resume source. Checkpoint writes only to the destination configured when the server starts and atomically replaces that file. It preserves controller state for any runnable model implementing the `SimulationController` protocol, including the legacy compatibility adapter. + +## Stop and restart + +Stop ends this server process and releases its listening port. It is authenticated through the same token and exact-origin WebSocket upgrade as every other command. The reader handles Stop immediately, including while an earlier command on that same socket is executing. Ordinary commands retain per-client ordering in a bounded queue of 32; additional queued commands receive an error rather than blocking Stop. + +Shutdown is cooperative: an in-progress individual simulation step finishes, then the rest of its batch is skipped. An already-running reset, scene capture, or atomic checkpoint write also finishes. There is no timeout that kills the worker in the middle of model state mutation or file replacement. A model operation that never returns will therefore keep the session in `stopping`. Queued operations and newly received Frame, Play, Step, Pause, Reset, and Checkpoint commands are rejected once stopping begins; repeated Stop requests are idempotent. Closing the final browser connection still pauses the session and allows reconnection. + +Browser Stop and terminal Ctrl+C share the same worker/socket cleanup path. After `stopped`, the server closes client sockets and its application runner, exits, and releases the port. Launch the next `microsimulator view` command from the terminal and open its newly printed URL; each process has a new token. The old browser keeps its last rendered frame and displays Stopped. diff --git a/python/src/microsimulator/viewer_server.py b/python/src/microsimulator/viewer_server.py index b0d6bf2..01decb8 100644 --- a/python/src/microsimulator/viewer_server.py +++ b/python/src/microsimulator/viewer_server.py @@ -13,6 +13,7 @@ from dataclasses import dataclass from functools import partial from pathlib import Path +from threading import Event from typing import Any, Literal, cast from urllib.parse import urlencode @@ -24,9 +25,10 @@ MAX_COMMAND_BYTES = 4096 MAX_STEP_BATCH = 10_000 +MAX_QUEUED_COMMANDS = 32 type ModelFactory = Callable[[], tuple[RunnableModel, Mapping[str, JSONValue]]] -type CommandName = Literal["frame", "step", "play", "pause", "reset", "checkpoint"] +type CommandName = Literal["frame", "step", "play", "pause", "reset", "checkpoint", "stop"] class LiveViewerError(RuntimeError): @@ -64,7 +66,7 @@ def parse_command(encoded: str) -> LiveCommand: ): raise LiveViewerError(f"step count must be an integer in [1, {MAX_STEP_BATCH}]") return LiveCommand("step", steps) - names: set[str] = {"frame", "play", "pause", "reset", "checkpoint"} + names: set[str] = {"frame", "play", "pause", "reset", "checkpoint", "stop"} if not isinstance(name, str) or name not in names: raise LiveViewerError("unknown command type") if set(command) != {"type"}: @@ -92,6 +94,11 @@ def __init__( self._model, self._provenance = self._build() self._completed_steps = 0 self._revision = 0 + self._stop_requested = Event() + + def request_stop(self) -> None: + """Signal the worker without waiting for its current operation.""" + self._stop_requested.set() @property def completed_steps(self) -> int: @@ -112,6 +119,8 @@ def step(self, steps: int = 1) -> None: completed = 0 try: for _ in range(steps): + if self._stop_requested.is_set(): + break self._model.step(self._dt) completed += 1 self._completed_steps += 1 @@ -174,11 +183,30 @@ def __init__(self, session: LiveSession, *, frame_steps: int = 1, fps: float = 3 self._play_wakeup = asyncio.Event() self._worker = ThreadPoolExecutor(max_workers=1, thread_name_prefix="microsimulator-live") self._operation_lock = asyncio.Lock() + self._close_task: asyncio.Task[None] | None = None + self.stopped = asyncio.Event() + + @property + def stopping(self) -> bool: + return self._close_task is not None + + def _require_active(self) -> None: + if self.stopping: + raise LiveViewerError("live session is stopping or stopped") async def _run(self, operation: Callable[..., Any], *arguments: object) -> Any: async with self._operation_lock: + self._require_active() loop = asyncio.get_running_loop() - return await loop.run_in_executor(self._worker, operation, *arguments) + work = loop.run_in_executor(self._worker, operation, *arguments) + try: + return await asyncio.shield(work) + except asyncio.CancelledError: + # Canceling an asyncio waiter does not cancel native/thread work. + # Keep the operation lock until the worker has really finished. + with suppress(Exception): + await asyncio.shield(work) + raise async def _message(self) -> dict[str, JSONValue]: return cast( @@ -199,7 +227,11 @@ async def broadcast_frame(self) -> None: stale.append(socket) continue try: - await socket.send_str(encoded) + await asyncio.wait_for(socket.send_str(encoded), timeout=1.0) + except TimeoutError: + # Retain this socket for shutdown even if it cannot drain a + # frame; Stop must not wait indefinitely on frame delivery. + continue except (ConnectionError, RuntimeError): stale.append(socket) self._sockets.difference_update(stale) @@ -207,20 +239,26 @@ async def broadcast_frame(self) -> None: async def _play(self) -> None: loop = asyncio.get_running_loop() try: - while self.playing and self._sockets: + while self.playing and self._sockets and not self.stopping: started = loop.time() await self._run(self.session.step, self.frame_steps) + if self.stopping: + break await self.broadcast_frame() delay = self.frame_interval - (loop.time() - started) if delay > 0.0: with suppress(TimeoutError): await asyncio.wait_for(self._play_wakeup.wait(), timeout=delay) self._play_wakeup.clear() + except Exception as error: + if not self.stopping: + await self._broadcast({"type": "error", "message": str(error)}) finally: self.playing = False self._play_task = None async def play(self) -> None: + self._require_active() if self.playing: return self.playing = True @@ -233,11 +271,15 @@ async def pause(self, *, broadcast: bool = True) -> None: self._play_wakeup.set() task = self._play_task if task is not None and task is not asyncio.current_task(): - await task - if broadcast: + await asyncio.shield(task) + if broadcast and not self.stopping: await self.broadcast_frame() async def command(self, command: LiveCommand) -> str | None: + if command.name == "stop": + self.request_stop() + return None + self._require_active() if command.name == "frame": await self.broadcast_frame() elif command.name == "step": @@ -258,21 +300,51 @@ async def command(self, command: LiveCommand) -> str | None: return None async def connect(self, socket: web.WebSocketResponse) -> None: + self._require_active() self._sockets.add(socket) await self._send_frame(socket) async def disconnect(self, socket: web.WebSocketResponse) -> None: self._sockets.discard(socket) - if not self._sockets: + if not self._sockets and not self.stopping: await self.pause(broadcast=False) - async def close(self) -> None: + async def _broadcast(self, message: dict[str, JSONValue]) -> None: + for socket in tuple(self._sockets): + if not socket.closed: + with suppress(ConnectionError, RuntimeError, TimeoutError): + # A stalled browser must not hold the shutdown admission + # path open indefinitely while its send buffer is full. + await asyncio.wait_for(socket.send_json(message), timeout=1.0) + + def request_stop(self) -> None: + """Begin idempotent shutdown outside any socket/command task.""" + if self.stopping: + return + self.session.request_stop() + self.playing = False + self._play_wakeup.set() + self._close_task = asyncio.create_task(self._close(), name="microsimulator-live-stop") + + async def _close(self) -> None: + await self._broadcast({"type": "session", "state": "stopping"}) await self.pause(broadcast=False) + # Wait for a manual batch, reset, frame capture, or atomic checkpoint. + # New/queued operations fail the admission check inside this lock. + async with self._operation_lock: + await asyncio.to_thread(self._worker.shutdown, wait=True, cancel_futures=True) + await self._broadcast({"type": "session", "state": "stopped"}) sockets = tuple(self._sockets) self._sockets.clear() for socket in sockets: - await socket.close(code=1001, message=b"server shutdown") - self._worker.shutdown(wait=True, cancel_futures=True) + with suppress(ConnectionError, RuntimeError): + await socket.close(code=1000, message=b"session stopped") + self.stopped.set() + + async def close(self) -> None: + self.request_stop() + assert self._close_task is not None + await asyncio.shield(self._close_task) _CONTROLLER_KEY = web.AppKey("microsimulator.controller", LiveController) @@ -290,11 +362,34 @@ def _authorized(request: web.Request) -> bool: async def _websocket(request: web.Request) -> web.StreamResponse: if not _authorized(request): raise web.HTTPForbidden(text="invalid live-viewer authority") + controller = request.app[_CONTROLLER_KEY] + if controller.stopping: + raise web.HTTPServiceUnavailable(text="live session is stopping or stopped") socket = web.WebSocketResponse(max_msg_size=MAX_COMMAND_BYTES, heartbeat=20.0) await socket.prepare(request) - controller = request.app[_CONTROLLER_KEY] - await controller.connect(socket) + commands: asyncio.Queue[LiveCommand] = asyncio.Queue(maxsize=MAX_QUEUED_COMMANDS) + + async def send_error(error: Exception) -> None: + if not socket.closed: + with suppress(ConnectionError, RuntimeError): + await socket.send_json({"type": "error", "message": str(error)}) + + async def execute_commands() -> None: + while True: + command = await commands.get() + try: + result = await controller.command(command) + if result is not None and not socket.closed: + await socket.send_json({"type": "checkpoint", "path": result}) + except Exception as error: + await send_error(error) + + # Read continuously so Stop on this socket can interrupt its own long batch. + # A bounded per-client queue retains ordinary command order without creating + # an unbounded number of tasks or blocking Stop behind queue backpressure. + consumer = asyncio.create_task(execute_commands(), name="microsimulator-live-commands") try: + await controller.connect(socket) async for message in socket: if message.type is not WSMsgType.TEXT: if message.type is WSMsgType.ERROR: @@ -302,18 +397,21 @@ async def _websocket(request: web.Request) -> web.StreamResponse: await socket.send_json({"type": "error", "message": "text commands required"}) continue try: - result = await controller.command(parse_command(cast(str, message.data))) - if result is not None: - await socket.send_json( - {"type": "checkpoint", "path": result}, - dumps=lambda value: json.dumps(value, separators=(",", ":")), - ) + command = parse_command(cast(str, message.data)) + if command.name == "stop": + controller.request_stop() + else: + if controller.stopping: + raise LiveViewerError("live session is stopping or stopped") + if commands.full(): + raise LiveViewerError("too many queued commands") + commands.put_nowait(command) except Exception as error: - await socket.send_json( - {"type": "error", "message": str(error)}, - dumps=lambda value: json.dumps(value, separators=(",", ":")), - ) + await send_error(error) finally: + consumer.cancel() + with suppress(asyncio.CancelledError): + await consumer await controller.disconnect(socket) return socket @@ -355,7 +453,8 @@ async def cleanup(_: web.Application) -> None: application.router.add_get("/", index_response) application.router.add_get("/api/v1/session", _websocket) application.router.add_static("/assets", assets, show_index=False) - application.on_cleanup.append(cleanup) + # Stop before aiohttp waits for active WebSocket request handlers. + application.on_shutdown.append(cleanup) return application, authority @@ -369,7 +468,7 @@ def serve_live( fps: float = 30.0, open_browser: bool = False, ) -> None: - """Serve one live session on loopback until interrupted.""" + """Serve one live session on loopback until Stop or terminal interruption.""" if host not in {"127.0.0.1", "::1", "localhost"}: raise LiveViewerError("live viewer host must be a loopback address") @@ -398,7 +497,7 @@ async def run() -> None: print(f"MicroSimulator live viewer: {url}", flush=True) if open_browser: webbrowser.open(url) - await asyncio.Event().wait() + await application[_CONTROLLER_KEY].stopped.wait() finally: await runner.cleanup() diff --git a/python/tests/test_viewer_shutdown.py b/python/tests/test_viewer_shutdown.py new file mode 100644 index 0000000..dd0c13c --- /dev/null +++ b/python/tests/test_viewer_shutdown.py @@ -0,0 +1,336 @@ +"""Shutdown regressions use real sockets and a controllably blocked worker.""" + +from __future__ import annotations + +import asyncio +import signal +import socket +import sys +from pathlib import Path +from threading import Event +from threading import enumerate as threads +from typing import Any, cast +from urllib.parse import urlsplit + +import pytest +from aiohttp import ClientSession, WSMsgType, web +from aiohttp.test_utils import TestClient, TestServer +from microsimulator import BackendKind, CellInit, Simulation, load_checkpoint +from microsimulator import checkpoint as checkpoint_module +from microsimulator.checkpoint import JSONValue +from microsimulator.viewer_server import ( + LiveCommand, + LiveController, + LiveSession, + LiveViewerError, + create_live_app, + parse_command, +) + + +def _factory() -> tuple[Simulation, dict[str, JSONValue]]: + simulation = Simulation(BackendKind.CPU) + simulation.add_cell(CellInit()) + return simulation, {} + + +def _dist(path: Path) -> Path: + dist = path / "dist" + (dist / "assets").mkdir(parents=True) + (dist / "index.html").write_text("shutdown test") + return dist + + +class _BlockedModel: + def __init__(self) -> None: + self.simulation, _ = _factory() + self.entered = Event() + self.release = Event() + + def controller_state(self) -> JSONValue: + return None + + def step(self, dt: float) -> None: + self.entered.set() + assert self.release.wait(10), "test did not release blocked step" + self.simulation.step(dt) + + +async def _entered(event: Event) -> None: + assert await asyncio.to_thread(event.wait, 5), "worker never reached blocking point" + + +@pytest.mark.parametrize("playing", [False, True]) +def test_stop_preempts_same_socket_batch_and_rejects_new_commands( + tmp_path: Path, + playing: bool, +) -> None: + async def exercise() -> None: + model = _BlockedModel() + session = LiveSession(lambda: (model, {}), dt=0.1) + app, token = create_live_app(session, _dist(tmp_path), frame_steps=10_000) + client = TestClient(TestServer(app)) + try: + await client.start_server() + origin = str(client.make_url("/")).rstrip("/") + ws = await client.ws_connect( + f"/api/v1/session?token={token}", + headers={"Origin": origin}, + ) + await ws.receive_json(timeout=5) + await ws.send_json({"type": "play"} if playing else {"type": "step", "steps": 10_000}) + await _entered(model.entered) + await ws.send_json({"type": "stop"}) + await ws.send_json({"type": "stop"}) + # Read until admission is closed while the current step is blocked. + while (await ws.receive_json(timeout=5)).get("type") != "session": + pass + for name in ("play", "step", "reset", "checkpoint"): + await ws.send_json({"type": name}) + rejected = await ws.receive_json(timeout=5) + assert rejected["type"] == "error" + assert "stopping" in rejected["message"] + assert session.completed_steps == 0 + model.release.set() + while True: + message = await ws.receive_json(timeout=5) + if message == {"type": "session", "state": "stopped"}: + break + assert session.completed_steps == 1 + assert (await ws.receive(timeout=5)).type == WSMsgType.CLOSE + assert ws.close_code == 1000 + finally: + model.release.set() + await client.close() + assert not any(t.name.startswith("microsimulator-live") for t in threads()) + + asyncio.run(exercise()) + + +def test_stop_waits_for_atomic_checkpoint_replace( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + async def exercise() -> None: + output = tmp_path / "run.json" + session = LiveSession(_factory, dt=0.1, checkpoint_output=output) + session.checkpoint() + original = output.read_bytes() + session.step() + entered, release = Event(), Event() + replace = checkpoint_module.os.replace + + def blocked_replace(source: str | Path, destination: str | Path) -> None: + entered.set() + assert release.wait(10), "test did not release checkpoint replace" + replace(source, destination) + + monkeypatch.setattr(checkpoint_module.os, "replace", blocked_replace) + controller = LiveController(session) + save = asyncio.create_task(controller.command(LiveCommand("checkpoint"))) + try: + await _entered(entered) + await controller.command(LiveCommand("stop")) + assert not controller.stopped.is_set() + assert output.read_bytes() == original + for name in ("play", "step", "reset", "checkpoint"): + with pytest.raises(LiveViewerError, match="stopping"): + await controller.command(parse_command('{"type":"' + name + '"}')) + release.set() + assert await save == str(output.resolve()) + await asyncio.wait_for(controller.close(), 5) + assert abs(load_checkpoint(output).time - 0.1) < 1e-7 + assert list(tmp_path.iterdir()) == [output] + finally: + release.set() + await controller.close() + + asyncio.run(exercise()) + + +def test_cancelled_waiter_does_not_release_worker_early() -> None: + async def exercise() -> None: + model = _BlockedModel() + session = LiveSession(lambda: (model, {}), dt=0.1) + controller = LiveController(session) + operation = asyncio.create_task(controller.command(LiveCommand("step", 10_000))) + try: + await _entered(model.entered) + operation.cancel() + controller.request_stop() + await asyncio.sleep(0) + assert not operation.done() + assert not controller.stopped.is_set() + model.release.set() + with pytest.raises(asyncio.CancelledError): + await operation + await asyncio.wait_for(controller.close(), 5) + assert session.completed_steps == 1 + finally: + model.release.set() + await controller.close() + + asyncio.run(exercise()) + + +def test_disconnect_pauses_without_stopping_and_close_is_idempotent(tmp_path: Path) -> None: + async def exercise() -> None: + session = LiveSession(_factory, dt=0.1) + app, token = create_live_app(session, _dist(tmp_path)) + client = TestClient(TestServer(app)) + await client.start_server() + origin = str(client.make_url("/")).rstrip("/") + try: + ws = await client.ws_connect( + f"/api/v1/session?token={token}", + headers={"Origin": origin}, + ) + await ws.receive_json(timeout=5) + await ws.send_json({"type": "play"}) + await ws.receive_json(timeout=5) + await ws.close() + ws2 = await client.ws_connect( + f"/api/v1/session?token={token}", + headers={"Origin": origin}, + ) + assert (await ws2.receive_json(timeout=5))["playing"] is False + await ws2.send_json({"type": "stop"}) + await asyncio.gather(ws2.close(), client.close()) + await client.close() + finally: + await client.close() + + asyncio.run(exercise()) + + +def test_stop_command_is_closed() -> None: + assert parse_command('{"type":"stop"}') == LiveCommand("stop") + with pytest.raises(LiveViewerError, match="unknown fields"): + parse_command('{"type":"stop","force":true}') + + +def test_stop_is_not_blocked_by_a_stalled_frame_send( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + async def exercise() -> None: + entered = asyncio.Event() + send = web.WebSocketResponse.send_str + frame_count = 0 + + async def stalled_send( + ws: web.WebSocketResponse, + data: str, + compress: int | None = None, + ) -> None: + nonlocal frame_count + if data.startswith('{"type":"frame"'): + frame_count += 1 + if frame_count > 1: + entered.set() + await asyncio.Event().wait() + await send(ws, data, compress=compress) + + monkeypatch.setattr(web.WebSocketResponse, "send_str", stalled_send) + app, token = create_live_app(LiveSession(_factory, dt=0.1), _dist(tmp_path)) + client = TestClient(TestServer(app)) + try: + await client.start_server() + origin = str(client.make_url("/")).rstrip("/") + ws = await client.ws_connect( + f"/api/v1/session?token={token}", + headers={"Origin": origin}, + ) + await ws.receive_json(timeout=5) + await ws.send_json({"type": "play"}) + await asyncio.wait_for(entered.wait(), 5) + await ws.send_json({"type": "stop"}) + assert await ws.receive_json(timeout=5) == {"type": "session", "state": "stopping"} + assert await ws.receive_json(timeout=5) == {"type": "session", "state": "stopped"} + await ws.close() + finally: + await client.close() + + asyncio.run(exercise()) + + +@pytest.mark.parametrize("termination", ["stop", "interrupt"]) +def test_cli_process_exits_and_reuses_port_for_another_model( + tmp_path: Path, + termination: str, +) -> None: + if termination == "interrupt" and sys.platform == "win32": + pytest.skip("Windows Ctrl+C requires an attached console; use documented manual check") + + async def exercise() -> None: + dist = _dist(tmp_path) + with socket.socket() as reservation: + reservation.bind(("127.0.0.1", 0)) + port = reservation.getsockname()[1] + for index in range(3): + model = tmp_path / f"model {index}.py" + model.write_text( + "from microsimulator import CellInit\n" + "def build(context):\n" + " simulation = context.simulation()\n" + " cell = CellInit()\n" + f" cell.length = {2 + index}\n" + " simulation.add_cell(cell)\n" + " return simulation\n" + ) + process = await asyncio.create_subprocess_exec( + sys.executable, + "-m", + "microsimulator", + "view", + "--model", + str(model), + "--backend", + "cpu", + "--dt", + "0.1", + "--viewer-dist", + str(dist), + "--port", + str(port), + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + try: + assert process.stdout is not None + line = (await asyncio.wait_for(process.stdout.readline(), 10)).decode().strip() + if not line.startswith("MicroSimulator live viewer: "): + _, error = await asyncio.wait_for(process.communicate(), 5) + pytest.fail(f"viewer did not start: {line} {error.decode()}") + url = urlsplit(line.split(": ", 1)[1]) + async with ( + ClientSession() as client, + client.ws_connect( + f"http://{url.netloc}/api/v1/session?{url.query}", + headers={"Origin": f"http://{url.netloc}"}, + ) as ws, + ): + frame = cast(dict[str, Any], await ws.receive_json(timeout=5)) + assert frame["scene"]["frame"]["cells"][0]["length"] == 2 + index + if termination == "interrupt": + process.send_signal(signal.SIGINT) + else: + await ws.send_json({"type": "stop"}) + assert await ws.receive_json(timeout=5) == { + "type": "session", + "state": "stopping", + } + assert await ws.receive_json(timeout=5) == { + "type": "session", + "state": "stopped", + } + assert (await ws.receive(timeout=5)).type == WSMsgType.CLOSE + _, error = await asyncio.wait_for(process.communicate(), 10) + assert process.returncode == 0, error.decode() + assert not error, error.decode() + finally: + if process.returncode is None: + process.kill() + await process.wait() + + asyncio.run(exercise()) diff --git a/viewer/README.md b/viewer/README.md index 1eea13e..75074fd 100644 --- a/viewer/README.md +++ b/viewer/README.md @@ -33,7 +33,28 @@ uv run microsimulator view \ --open ``` -Without `--open`, open the tokenized loopback URL printed by `microsimulator`. The live transport can play, pause, advance one step, rebuild the original model, and write to the configured checkpoint destination. Camera position, display mapping, grid slice, and selected-cell identity survive frame updates. +Without `--open`, open the tokenized loopback URL printed by `microsimulator`. The live transport can play, pause, advance one step, rebuild the original model, write to the configured checkpoint destination, and stop the session. Camera position, display mapping, grid slice, and selected-cell identity survive frame updates. + +### Stop one model and start another + +Click **Stop session** or press **Ctrl+C** once in the terminal running the server. The current individual step or checkpoint write finishes, the browser displays **Stopped**, and the command returns to the prompt. A large playback batch does not have to finish. Start another `microsimulator view` command using the same port and open the new printed URL. Reset rebuilds the current model; Pause keeps its process available; closing the browser pauses it and allows reconnection. Stop does not automatically save a checkpoint: use Checkpoint first if you need restartable state. + +For example, after stopping the trap model above, launch a different model on the same default port: + +```console +uv run microsimulator view --model examples/tutorials/biophysics.py --backend cpu --seed 42 --dt 0.02 --port 8765 --open +``` + +The command is a single line and also works in PowerShell where the Python/native build is available. To distinguish Windows console behavior from browser behavior, use this manual verification procedure in an attached PowerShell or Command Prompt console: + +1. Record the Windows version, terminal application/version, Python version, and exact launch command. Start the command above and click Stop while paused. Confirm the prompt returns, then start the second model on port 8765. +2. Repeat with Play active and `--frame-steps 10000`. Confirm Stopping transitions to Stopped without finishing the entire batch. +3. Repeat using Ctrl+C once, both paused and playing. Confirm the prompt returns without `taskkill`, then immediately start another model on the same port. +4. Close only the browser tab during Play, then reopen the printed URL. Confirm the session remains available and paused. + +The automated `python/tests/test_viewer_shutdown.py` suite covers same-socket Stop, checkpoint completion, worker cleanup, and repeated real subprocess restarts. Its Stop/port-reuse test is portable to Windows; the SIGINT subprocess case runs on POSIX. Windows console Ctrl+C must be checked with the attached-console procedure above; passing the POSIX case does not establish Windows behavior. + +The `Windows live-session shutdown` GitHub Actions job builds the CPU extension on `windows-2025` and runs the portable server and shutdown tests with dependencies from `uv.lock`. Its uploaded report records Windows, PowerShell, Python, backend availability, and individual test results. The detached CI process cannot substitute for the attached-console Ctrl+C check. ## Capabilities diff --git a/viewer/browser/live-shutdown.mjs b/viewer/browser/live-shutdown.mjs new file mode 100644 index 0000000..13a6ed6 --- /dev/null +++ b/viewer/browser/live-shutdown.mjs @@ -0,0 +1,174 @@ +// Real Python server + browser, including process exit and same-port restart. +// Run from the repository root after building this worktree's Python/viewer. +import assert from "node:assert/strict"; +import { spawn } from "node:child_process"; +import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import { createInterface } from "node:readline"; + +const { chromium, expect } = await import( + process.env.MICROSIMULATOR_PLAYWRIGHT_MODULE ?? "@playwright/test" +); +const evidence = + process.env.EVIDENCE_DIR ?? "/tmp/microsimulator-live-shutdown"; +const python = + process.env.MICROSIMULATOR_PYTHON ?? path.resolve(".venv/bin/python"); +const port = process.env.VIEWER_PORT ?? "4327"; +const temporary = await mkdtemp(path.join(tmpdir(), "microsimulator-stop-")); +await mkdir(evidence, { recursive: true }); +await writeFile( + path.join(temporary, "model.py"), + ` +import time +from microsimulator import CellInit +class Model: + def __init__(self, context): + self.simulation = context.simulation() + self.simulation.add_cell(CellInit()) + def step(self, dt): + time.sleep(0.02) + self.simulation.step(dt) + def controller_state(self): + return {"kind": "shutdown-browser-fixture"} +def build(context): + return Model(context) +`, +); + +const browser = await chromium.launch({ headless: true }); +let processUnderTest; +const errors = []; +const results = []; +try { + for (const mode of ["paused", "playing", "reconnect", "interrupt"]) { + if (mode === "interrupt" && process.platform === "win32") continue; + const child = spawn( + python, + [ + "-m", + "microsimulator", + "view", + "--model", + path.join(temporary, "model.py"), + "--backend", + "cpu", + "--dt", + "0.01", + "--frame-steps", + "10000", + "--viewer-dist", + path.resolve("viewer/dist"), + "--port", + port, + "--checkpoint-output", + path.join(temporary, "checkpoint.json"), + ], + { stdio: ["ignore", "pipe", "pipe"] }, + ); + processUnderTest = child; + let stderr = ""; + child.stderr.on("data", (data) => { + stderr += data.toString(); + }); + const exit = new Promise((resolve) => + child.once("exit", (code, signal) => resolve({ code, signal })), + ); + const lines = createInterface({ input: child.stdout }); + const url = await new Promise((resolve, reject) => { + const timer = setTimeout( + () => reject(new Error("server startup timed out")), + 10000, + ); + lines.on("line", (line) => { + if (line.startsWith("MicroSimulator live viewer: ")) { + clearTimeout(timer); + resolve(line.slice("MicroSimulator live viewer: ".length)); + } + }); + child.once("exit", () => { + clearTimeout(timer); + reject(new Error(stderr)); + }); + child.once("error", reject); + }); + let page = await browser.newPage({ + viewport: { width: 1440, height: 960 }, + }); + page.on("pageerror", (error) => errors.push(error.message)); + await page.goto(url); + await expect(page.locator("#live-label")).toHaveText("Paused"); + await expect(page.locator("#live-stop")).toBeEnabled(); + if (mode === "reconnect") { + await page.close(); + // The existing viewer deliberately requires a minimum width of 880px. + page = await browser.newPage({ viewport: { width: 880, height: 844 } }); + page.on("pageerror", (error) => errors.push(error.message)); + await page.goto(url); + await expect(page.locator("#live-label")).toHaveText("Paused"); + assert.equal(child.exitCode, null); + await page.locator("#live-checkpoint").click(); + await expect(page.locator("#status")).toContainText("Checkpoint saved"); + assert.ok( + (await readFile(path.join(temporary, "checkpoint.json"))).length > 100, + ); + } + if (mode === "playing") { + await page.locator("#live-play").click(); + await expect(page.locator("#live-label")).toHaveText("Running"); + } + const start = performance.now(); + if (mode === "interrupt") { + child.kill("SIGINT"); + } else { + await page.locator("#live-stop").focus(); + await page.keyboard.press("Enter"); + } + await expect(page.locator("#live-label")).toHaveText("Stopped", { + timeout: 5000, + }); + const toolbar = await page.locator("#live-transport").boundingBox(); + const viewport = await page.locator("#canvas-host").boundingBox(); + assert.ok(toolbar.x >= viewport.x); + assert.ok(toolbar.x + toolbar.width <= viewport.x + viewport.width); + await expect(page.locator("#status")).toContainText("Session stopped"); + for (const id of ["play", "step", "reset", "checkpoint", "stop"]) { + await expect(page.locator(`#live-${id}`)).toBeDisabled(); + } + const stopped = await Promise.race([ + exit, + new Promise((_, reject) => { + const timer = setTimeout( + () => reject(new Error("server did not exit")), + 5000, + ); + timer.unref(); + }), + ]); + assert.deepEqual(stopped, { code: 0, signal: null }); + assert.equal(stderr, ""); + await page.screenshot({ path: `${evidence}/${mode}-stopped.png` }); + results.push({ + mode, + stopMilliseconds: performance.now() - start, + exitCode: stopped.code, + }); + await page.close(); + lines.close(); + processUnderTest = undefined; + } + assert.deepEqual(errors, []); + const result = { + result: "passed", + browser: browser.version(), + platform: process.platform, + port, + results, + }; + await writeFile(`${evidence}/results.json`, JSON.stringify(result, null, 2)); + console.log(JSON.stringify(result, null, 2)); +} finally { + processUnderTest?.kill("SIGKILL"); + await browser.close(); + await rm(temporary, { recursive: true, force: true }); +} diff --git a/viewer/index.html b/viewer/index.html index 0b1a8a3..7a2945d 100644 --- a/viewer/index.html +++ b/viewer/index.html @@ -191,6 +191,14 @@

Inspect a colony snapshot

> Checkpoint +