diff --git a/.github/workflows/live-shutdown-windows.yml b/.github/workflows/live-shutdown-windows.yml new file mode 100644 index 0000000..5fee6e0 --- /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 Stop, console Ctrl+C, worker draining, and same-port 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/assets/capsules/colony-after.png b/docs/assets/capsules/colony-after.png new file mode 100644 index 0000000..eaf347a Binary files /dev/null and b/docs/assets/capsules/colony-after.png differ diff --git a/docs/assets/capsules/colony-before.png b/docs/assets/capsules/colony-before.png new file mode 100644 index 0000000..71fa385 Binary files /dev/null and b/docs/assets/capsules/colony-before.png differ diff --git a/docs/assets/capsules/isolated-after.png b/docs/assets/capsules/isolated-after.png new file mode 100644 index 0000000..374b885 Binary files /dev/null and b/docs/assets/capsules/isolated-after.png differ diff --git a/docs/assets/capsules/isolated-before.png b/docs/assets/capsules/isolated-before.png new file mode 100644 index 0000000..5ea415d Binary files /dev/null and b/docs/assets/capsules/isolated-before.png differ diff --git a/docs/capsule-rendering-validation.md b/docs/capsule-rendering-validation.md new file mode 100644 index 0000000..95c211a --- /dev/null +++ b/docs/capsule-rendering-validation.md @@ -0,0 +1,40 @@ +# Capsule rendering validation + +The previous renderer (`b69193b9`) combined a closed cylinder with two complete spheres. The sphere tessellation remained in world orientation while the cylinder followed the cell direction. Their polygonal boundaries did not coincide, and the cylinder end disks intersected the spherical surfaces. The resulting narrow lines at the cap joins are visible in the baseline screenshots below. + +The replacement uses an open cylinder and two hemispheres, all sharing the cell orientation. Both hemispherical equators use the cylinder's exact radial samples and normals. Radius scales the caps uniformly; cylindrical length scales only the cylinder and displaces cap centers. The selection wireframe uses these same surfaces with an 8% radius expansion. Zero centerline length collapses the cylinder and joins two hemispheres into a sphere. + +## Identical-camera comparisons + +These are unmodified screenshots from `viewer/browser/capsules.mjs`, using the same geometry fixtures, camera, lights, colors, browser, and viewport. The baseline renderer was recorded before the implementation changed. Neither image comes from a different simulation trajectory. + +| Fixture | Before | After | +| --- | --- | --- | +| Isolated rod, arbitrary 3D direction | ![Baseline isolated rod with a visible cap join](assets/capsules/isolated-before.png) | ![Continuous isolated capsule](assets/capsules/isolated-after.png) | +| Dense colony, near view | ![Baseline cap rings throughout the colony](assets/capsules/colony-before.png) | ![Colony without the cap rings](assets/capsules/colony-after.png) | + +The cap-join lines disappear at the same camera positions in the near and distant fixtures. The browser harness also records 48 frames of prescribed changes in position and orientation, so the surface can be inspected during movement. Silhouette tessellation, pixel aliasing, and real motion remain possible; this change does not smooth simulation state or claim to eliminate every source of shimmer. + +## Geometry and interaction checks + +The Vitest geometry suite checks matched seam positions/normals, outward-facing hemisphere triangles, absence of cylinder end disks, constant distance from the centerline segment, exact axial extents, conservative bounds, arbitrary directions, and zero-length cells. It also ray-picks both capsule tips with real Three.js instanced meshes and checks selection-radius inflation. The browser harness exercises color attributes on every mesh, pointer picking, stable-ID selection through reordered frames and removal, actual selected-overlay surface positions, and a selected zero-length sphere. The largest measured overlay-surface error in the browser fixture was `1.29e-8` world units. + +## Local rendering and resource measurements + +Both comparisons ran on macOS in headless Chromium `153.0.8010.12`, with the same browser launch configuration. The corrected run reports ANGLE/Vulkan SwiftShader: these are **software WebGL measurements**, not physical GPU performance results. The fixture uses 512 cells, ten warmup frames, and 60 measured complete frame replacements, including transforms, coloring, rendering, and `gl.finish()` synchronization. Small timing differences are within local measurement noise. + +| Measurement | Before | After | +| --- | ---: | ---: | +| Instanced draw calls for the colony | 3 | 3 | +| Shared geometry vertices (cylinder + both caps) | 388 | 400 | +| Triangles per cell | 504 | 576 | +| Triangles for 512 cells | 258,048 | 294,912 | +| Median replacement/render time | 2.2 ms | 2.1 ms | +| p95 replacement/render time | 2.4 ms | 2.4 ms | +| Renderer geometry count at both sampled frames | 4 | 4 | +| Tracked live WebGL buffers across 59 replacements | 390 → 744 | 18 → 18 | +| Tracked buffers after an empty frame | 732 | 0 | + +The geometry count alone concealed an existing resource leak: removing an instanced mesh and disposing its geometry did not release `instanceMatrix` and `instanceColor` buffers. Frame replacement now calls `InstancedMesh.dispose()` as well as disposing each shared geometry/material once. Browser instrumentation observes actual WebGL buffer creation/deletion; buffers remain bounded across replacements and return to zero after clearing the colony and disposing the viewer. + +The modest tessellation change (24 radial samples and six rows per hemisphere) improves silhouettes, but it is not the basis for the seam fix: the tests verify the changed surface topology and exact equator agreement. For reproduction commands, fixture details, videos, and resource assertions, see [the browser harness instructions](../viewer/browser/README.md). diff --git a/docs/protocols/live-viewer-v1.md b/docs/protocols/live-viewer-v1.md index fbbf780..98da8c8 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,19 @@ 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. + +Network delivery has separate deadlines from cooperative model work. Initial frame writes, broadcasts, and lifecycle notifications allow one second per receiver; independent lifecycle deliveries run concurrently. The complete WebSocket close operation, including writing and draining its close frame, also has a one-second deadline. An unresponsive connection is then aborted, discarding queued network bytes so neither its initial-send handler nor the application runner waits indefinitely for a reader. Responsive clients still receive `stopping`, `stopped`, and a normal code-1000 close. A stalled client may miss these notifications and observe an unexpected disconnect; this does not cancel an active simulation operation or interrupt an atomic checkpoint write. diff --git a/python/src/microsimulator/viewer_server.py b/python/src/microsimulator/viewer_server.py index 4e2599a..865ef29 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,12 @@ MAX_COMMAND_BYTES = 4096 MAX_STEP_BATCH = 10_000 +MAX_QUEUED_COMMANDS = 32 +SOCKET_SEND_TIMEOUT = 1.0 +SOCKET_CLOSE_TIMEOUT = 1.0 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 +68,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 +96,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: @@ -113,6 +122,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 @@ -179,15 +190,35 @@ def __init__(self, session: LiveSession, *, frame_steps: int = 1, fps: float = 3 self.frame_interval = 1.0 / fps self.playing = False self._sockets: set[web.WebSocketResponse] = set() + self._transports: dict[web.WebSocketResponse, asyncio.Transport | None] = {} self._play_task: asyncio.Task[None] | None = None 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( @@ -196,40 +227,66 @@ async def _message(self) -> dict[str, JSONValue]: ) async def _send_frame(self, socket: web.WebSocketResponse) -> None: - await socket.send_str(json.dumps(await self._message(), separators=(",", ":"))) + encoded = json.dumps(await self._message(), separators=(",", ":")) + await asyncio.wait_for(socket.send_str(encoded), timeout=SOCKET_SEND_TIMEOUT) async def broadcast_frame(self) -> None: - if not self._sockets: + sockets = tuple(self._sockets) + if not sockets: return encoded = json.dumps(await self._message(), separators=(",", ":")) stale: list[web.WebSocketResponse] = [] - for socket in tuple(self._sockets): + # A frame captured for earlier clients must not arrive ahead of a newly + # connected client's initial frame (or carry its stale playing flag). + for socket in sockets: + if self.stopping: + break if socket.closed: stale.append(socket) continue try: - await socket.send_str(encoded) + await asyncio.wait_for(socket.send_str(encoded), timeout=SOCKET_SEND_TIMEOUT) + 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) + except asyncio.CancelledError: + # Another writer can cancel aiohttp's shared drain waiter. + # An externally canceled task must still propagate cancellation. + task = asyncio.current_task() + if task is not None and task.cancelling(): + raise + transport = self._transports.get(socket) + if transport is not None: + transport.abort() + stale.append(socket) self._sockets.difference_update(stale) 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 @@ -242,11 +299,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": @@ -266,22 +327,100 @@ async def command(self, command: LiveCommand) -> str | None: return str(destination) return None - async def connect(self, socket: web.WebSocketResponse) -> None: + async def connect( + self, socket: web.WebSocketResponse, transport: asyncio.Transport | None = None + ) -> None: + self._require_active() + stale = {client for client in self._sockets if client.closed} + if stale: + self._sockets.difference_update(stale) + if not self._sockets: + # The peer may receive its close handshake before the old + # request handler reaches finally. Honor last-client pause + # before admitting a replacement connection in that window. + await self.pause(broadcast=False) + self._require_active() self._sockets.add(socket) + self._transports[socket] = transport await self._send_frame(socket) async def disconnect(self, socket: web.WebSocketResponse) -> None: self._sockets.discard(socket) - if not self._sockets: + self._transports.pop(socket, None) + 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: + async def send(socket: web.WebSocketResponse) -> None: + if not socket.closed: + try: + await asyncio.wait_for(socket.send_json(message), timeout=SOCKET_SEND_TIMEOUT) + except (ConnectionError, RuntimeError, TimeoutError): + pass + except asyncio.CancelledError: + task = asyncio.current_task() + if task is not None and task.cancelling(): + raise + transport = self._transports.get(socket) + if transport is not None: + transport.abort() + + # Each receiver gets the same bounded opportunity to consume a message; + # unresponsive clients do not add serial delays to session shutdown. + await asyncio.gather(*(send(socket) for socket in tuple(self._sockets))) + + 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) - sockets = tuple(self._sockets) + # 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((socket, self._transports.get(socket)) for socket in 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) + self._transports.clear() + await asyncio.gather(*(_close_socket(socket, transport) for socket, transport in sockets)) + self.stopped.set() + + async def close(self) -> None: + self.request_stop() + assert self._close_task is not None + await asyncio.shield(self._close_task) + + +async def _close_socket(socket: web.WebSocketResponse, transport: asyncio.Transport | None) -> None: + closed = False + try: + # aiohttp applies its own timeout only AFTER writing/draining the close + # frame. Bound the whole operation, including that preceding drain. + await asyncio.wait_for( + socket.close(code=1000, message=b"session stopped"), timeout=SOCKET_CLOSE_TIMEOUT + ) + closed = True + except (ConnectionError, RuntimeError, TimeoutError): + pass + except asyncio.CancelledError: + # aiohttp shares one drain waiter between writes. A timed-out initial + # send may cancel that waiter, independently of this cleanup task. + # Treat that as a failed close, but preserve real task cancellation. + task = asyncio.current_task() + if task is not None and task.cancelling(): + raise + finally: + if not closed and transport is not None: + # Transport.close() still tries to flush queued bytes; a receiver + # that never reads requires abort() to release the drain waiters. + transport.abort() _CONTROLLER_KEY = web.AppKey("microsimulator.controller", LiveController) @@ -299,11 +438,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, request.transport) async for message in socket: if message.type is not WSMsgType.TEXT: if message.type is WSMsgType.ERROR: @@ -311,19 +473,34 @@ 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) + except (ConnectionError, TimeoutError): + # Initial-frame delivery is also bounded; a stalled receiver has no + # authority to keep an upgraded request alive during runner cleanup. + pass finally: - await controller.disconnect(socket) + consumer.cancel() + try: + # Drop client authority and pause before waiting for canceled work + # to drain. Otherwise a reconnect during that wait keeps playback + # alive because the old socket still appears to be connected. + await controller.disconnect(socket) + finally: + try: + with suppress(asyncio.CancelledError): + await consumer + finally: + await _close_socket(socket, request.transport) return socket @@ -364,7 +541,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 @@ -378,7 +556,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") @@ -407,7 +585,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..8e8adcd --- /dev/null +++ b/python/tests/test_viewer_shutdown.py @@ -0,0 +1,646 @@ +"""Shutdown regressions use real sockets and a controllably blocked worker.""" + +from __future__ import annotations + +import asyncio +import os +import signal +import socket +import subprocess +import sys +from contextlib import suppress +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, ClientWebSocketResponse, WSMsgType, web +from aiohttp.test_utils import TestClient, TestServer +from microsimulator import BackendKind, CellInit, Simulation, Vec3, load_checkpoint +from microsimulator.checkpoint import JSONValue +from microsimulator.viewer_server import ( + _CONTROLLER_KEY, # pyright: ignore[reportPrivateUsage] + 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 = 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(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()) + + +@pytest.mark.parametrize("_attempt", range(10)) +def test_disconnect_pauses_without_stopping_and_close_is_idempotent( + tmp_path: Path, + _attempt: int, +) -> 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_reconnect_observes_pause_before_disconnected_command_drains( + tmp_path: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + async def exercise() -> None: + entered, release = Event(), Event() + reconnect_started = asyncio.Event() + frame_message = LiveSession.frame_message + connect = LiveController.connect + frames, connections = 0, 0 + + def blocked_frame(session: LiveSession, *, playing: bool) -> dict[str, JSONValue]: + nonlocal frames + frames += 1 + if frames == 2: + entered.set() + assert release.wait(10), "test did not release frame capture" + return frame_message(session, playing=playing) + + async def observe_connect( + controller: LiveController, + ws: web.WebSocketResponse, + transport: asyncio.Transport | None = None, + ) -> None: + nonlocal connections + connections += 1 + if connections == 2: + reconnect_started.set() + await connect(controller, ws, transport) + + monkeypatch.setattr(LiveSession, "frame_message", blocked_frame) + monkeypatch.setattr(LiveController, "connect", observe_connect) + 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("/") + + async def open_socket() -> ClientWebSocketResponse: + return await client.ws_connect( + f"/api/v1/session?token={token}", + headers={"Origin": origin}, + ) + + ws = await open_socket() + await ws.receive_json(timeout=5) + await ws.send_json({"type": "play"}) + await _entered(entered) + # close() finishes the network handshake while the server's canceled + # frame capture is still holding its worker/serialization lock. + await ws.close() + reconnect = asyncio.create_task(open_socket()) + await asyncio.wait_for(reconnect_started.wait(), 5) + release.set() + ws2 = await reconnect + assert (await ws2.receive_json(timeout=5))["playing"] is False + await ws2.close() + finally: + release.set() + await client.close() + + asyncio.run(exercise()) + + +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", "runner_cleanup"]) +def test_real_tcp_backpressure_releases_connections_and_reuses_port( + tmp_path: Path, termination: str +) -> None: + def large_factory() -> tuple[Simulation, dict[str, JSONValue]]: + simulation = Simulation(BackendKind.CPU) + cell = CellInit() + for index in range(40_000): + cell.position = Vec3(index * 3.0, 0.0, 0.0) + simulation.add_cell(cell) + return simulation, {} + + async def exercise() -> None: + dist = _dist(tmp_path) + app, token = create_live_app(LiveSession(large_factory, dt=0.1), dist) + controller = app[_CONTROLLER_KEY] + probe_transport: asyncio.Transport | None = None + probe_prepared = asyncio.Event() + + async def limit_probe_send_buffer( + request: web.Request, response: web.StreamResponse + ) -> None: + nonlocal probe_transport + if request.headers.get("X-Backpressure-Probe") == "1": + assert response.status == 101 + probe_transport = request.transport + assert probe_transport is not None + peer = cast(socket.socket, probe_transport.get_extra_info("socket")) + # Bound kernel buffering too: Windows overlapped writes may + # otherwise accept the whole scene before the peer consumes it. + peer.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, 16_384) + probe_prepared.set() + + app.on_response_prepare.append(limit_probe_send_buffer) + server = TestServer(app, host="127.0.0.1") + await server.start_server() + port = server.make_url("/").port + assert port is not None + origin = f"http://127.0.0.1:{port}" + stalled: asyncio.StreamWriter | None = None + cleanup: asyncio.Task[None] | None = None + receiver: asyncio.Task[None] | None = None + pump: asyncio.Task[None] | None = None + try: + async with ClientSession() as client: + healthy = await client.ws_connect( + f"{origin}/api/v1/session?token={token}", + headers={"Origin": origin}, + max_msg_size=128 * 1024 * 1024, + ) + initial = cast(dict[str, Any], await healthy.receive_json(timeout=15)) + assert len(initial["scene"]["frame"]["cells"]) == 40_000 + notices: asyncio.Queue[dict[str, Any]] = asyncio.Queue() + frames: asyncio.Queue[None] = asyncio.Queue() + + async def drain_healthy() -> None: + async for message in healthy: + assert message.type is WSMsgType.TEXT + body = cast(dict[str, Any], message.json()) + if body.get("type") == "frame": + frames.put_nowait(None) + else: + notices.put_nowait(body) + + receiver = asyncio.create_task(drain_healthy()) + # Limit the receive window before TCP negotiation, then stop + # reading before sending the authenticated WebSocket upgrade. + # No send/close implementation or timeout is mocked. + raw = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + try: + raw.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 4096) + raw.setblocking(False) + await asyncio.get_running_loop().sock_connect(raw, ("127.0.0.1", port)) + reader, stalled = await asyncio.open_connection(sock=raw) + except BaseException: + raw.close() + raise + cast(asyncio.Transport, stalled.transport).pause_reading() + stalled.write( + ( + f"GET /api/v1/session?token={token} HTTP/1.1\r\n" + f"Host: 127.0.0.1:{port}\r\nOrigin: {origin}\r\n" + "Upgrade: websocket\r\nConnection: Upgrade\r\n" + "Sec-WebSocket-Key: MTIzNDU2Nzg5MDEyMzQ1Ng==\r\n" + "Sec-WebSocket-Version: 13\r\nX-Backpressure-Probe: 1\r\n\r\n" + ).encode("ascii") + ) + await stalled.drain() + # Observe the server's authenticated 101 response without + # allowing the client to prefetch any scene payload first. + await asyncio.wait_for(probe_prepared.wait(), 5) + # Synchronize on observed backpressure instead of assuming that + # a delay was long enough for the initial send to fill the queue. + assert probe_transport is not None + blocked = probe_transport + requests_sent = 0 + + async def fill_windows_loopback() -> None: + nonlocal requests_sent + # Windows can accept the entire initial scene below the + # Proactor transport while the receiver is already paused. + # Exercise sustained real broadcasts on that platform; + # POSIX retains the stalled-initial-delivery scenario. + for _ in range(8): + if blocked.get_write_buffer_size() > 1_000_000: + return + await healthy.send_json({"type": "frame"}) + requests_sent += 1 + # The enclosing 15-second setup budget owns this wait; + # native capture is not part of the network deadline. + await frames.get() + + if sys.platform == "win32": + pump = asyncio.create_task(fill_windows_loopback()) + try: + async with asyncio.timeout(15): + while blocked.get_write_buffer_size() <= 1_000_000: + assert not blocked.is_closing(), ( + "probe closed before backpressure observed" + ) + if pump is not None and pump.done(): + await pump + raise AssertionError( + "8 real Frame requests did not cause backpressure" + ) + await asyncio.sleep(0.01) + except TimeoutError as error: + peer = cast(socket.socket, blocked.get_extra_info("socket")) + buffered = len(reader._buffer) # pyright: ignore[reportPrivateUsage] + chains: list[str] = [] + for task in list(asyncio.all_tasks())[:16]: + current = cast(Any, task.get_coro()) + chain: list[str] = [] + for _ in range(16): + if current is None: + break + code = getattr(current, "cr_code", None) + chain.append(code.co_name if code else type(current).__name__) + current = getattr(current, "cr_await", None) + chains.append(" -> ".join(chain)) + registered = len(controller._sockets) # pyright: ignore[reportPrivateUsage] + raise AssertionError( + f"no backpressure: transport={type(blocked).__name__}, " + f"queued={blocked.get_write_buffer_size()}, reader_bytes={buffered}, " + f"send_buffer={peer.getsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF)}, " + f"receive_buffer={raw.getsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF)}, " + f"registered={registered}, frame_requests={requests_sent}, " + f"await_chains={chains}" + ) from error + assert blocked.get_write_buffer_size() > 1_000_000 + if pump is not None: + # Only the test's request pump is canceled. Already-running + # server work is still drained by cooperative shutdown. + pump.cancel() + with suppress(asyncio.CancelledError): + await pump + if termination == "stop": + await healthy.send_json({"type": "stop"}) + else: + cleanup = asyncio.create_task(server.close()) + assert await asyncio.wait_for(notices.get(), 5) == { + "type": "session", + "state": "stopping", + } + assert await asyncio.wait_for(notices.get(), 5) == { + "type": "session", + "state": "stopped", + } + await asyncio.wait_for(receiver, 5) + assert healthy.close_code == 1000 + await asyncio.wait_for(controller.stopped.wait(), 5) + await asyncio.wait_for(cleanup or server.close(), 5) + # close() alone may still flush indefinitely; abort() releases + # queued bytes without resuming the paused receiver. + assert blocked.is_closing() + assert blocked.get_write_buffer_size() == 0 + assert not any(t.name.startswith("microsimulator-live") for t in threads()) + replacement, next_token = create_live_app(LiveSession(_factory, dt=0.1), dist) + next_server = TestServer(replacement, host="127.0.0.1", port=port) + await next_server.start_server() + try: + async with client.ws_connect( + f"{origin}/api/v1/session?token={next_token}", + headers={"Origin": origin}, + ) as next_ws: + frame = cast(dict[str, Any], await next_ws.receive_json(timeout=5)) + assert len(frame["scene"]["frame"]["cells"]) == 1 + finally: + await asyncio.wait_for(next_server.close(), 5) + finally: + tasks = [task for task in (pump, receiver) if task is not None] + for task in tasks: + task.cancel() + # Verification awaits failures above. Cleanup must still release + # real sockets if either helper had already failed. + await asyncio.gather(*tasks, return_exceptions=True) + if stalled is not None: + stalled.transport.abort() + await asyncio.wait_for(stalled.wait_closed(), 5) + if cleanup is not None: + await asyncio.wait_for(asyncio.shield(cleanup), 5) + await asyncio.wait_for(server.close(), 5) + + 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: + 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, + # Keep the test runner outside the console receiving Ctrl+C. + # CREATE_NEW_PROCESS_GROUP would disable Ctrl+C in the child. + creationflags=( + getattr(subprocess, "CREATE_NEW_CONSOLE", 0) + if termination == "interrupt" and sys.platform == "win32" + else 0 + ), + ) + 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": + if sys.platform == "win32": + # A separate sender attaches to the isolated viewer + # console. Ignore the event in the sender only; the + # viewer receives the real Windows CTRL_C_EVENT. + sender = await asyncio.create_subprocess_exec( + sys.executable, + "-c", + "import ctypes, sys\n" + "kernel = ctypes.WinDLL('kernel32', use_last_error=True)\n" + "kernel.FreeConsole()\n" + "for operation, arguments in (\n" + " (kernel.AttachConsole, (int(sys.argv[1]),)),\n" + " (kernel.SetConsoleCtrlHandler, (None, True)),\n" + " (kernel.GenerateConsoleCtrlEvent, (0, 0)),\n" + "):\n" + " if not operation(*arguments):\n" + " raise ctypes.WinError(ctypes.get_last_error())\n", + str(process.pid), + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + try: + _, sender_error = await asyncio.wait_for(sender.communicate(), 10) + assert sender.returncode == 0, sender_error.decode() + finally: + if sender.returncode is None: + sender.kill() + await sender.wait() + else: + 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 87ae74f..09b59ea 100644 --- a/viewer/README.md +++ b/viewer/README.md @@ -33,13 +33,34 @@ 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. It sends SIGINT on POSIX. On Windows it starts each viewer in an isolated console and uses a separate attached sender to deliver a real [Windows CTRL_C_EVENT](https://learn.microsoft.com/en-us/windows/console/generateconsolectrlevent), leaving the test runner unaffected. Both paths verify orderly browser notifications, clean process exit, and three different models reusing the same port. This exercises the operating-system interruption path; use the manual procedure above to check a particular interactive terminal application and keyboard configuration. + +The `Windows live-session shutdown` GitHub Actions job builds the CPU extension on `windows-2025` and runs the server and shutdown tests, including isolated-console Ctrl+C, with dependencies from `uv.lock`. Its uploaded report records Windows, PowerShell, Python, backend availability, and individual test results. ## Capabilities - SHA-256 verification over the Python writer's RFC 8785 canonical frame; - strict scene v2 structural and numerical validation; -- instanced cylinder and sphere rendering for exact spherocylinder geometry; +- instanced open cylinders and matching hemispheres for continuous capsule surfaces; - device walls rendered from plane, sphere, box, and cylinder constraints; - orbit, pan, zoom, colony framing, raycast picking, and selection highlighting; - a draggable camera-synchronized flat-corner view cube with readable labels, shortest-path single-click snapping, and double-click label leveling; @@ -63,6 +84,10 @@ pnpm --dir viewer build The unit suite includes a Python-authored scene fixture whose digest contains floating-point values that ordinary Python and JavaScript JSON serializers spell differently. Passing that test is the cross-language integrity gate. +Capsule geometry tests verify the scene's cylindrical centerline length, constant radius, spherical ends, zero-length sphere case, arbitrary orientation, outward topology, exact equator positions/normals, and ray picking. The selected-cell overlay uses the same geometry with radius increased by 8%; the cap centers retain the original centerline length. Three instanced draw calls represent the colony, independent of cell count. Replacing frames disposes both geometry and instance buffers. + +For visual and GPU-resource checks, see [the capsule browser regression](browser/README.md). Mesh tessellation and pixel aliasing can still affect silhouettes at distant zoom levels; the shared tangent joins specifically remove overlapping end disks and mismatched sphere/cylinder boundaries. The viewer does not smooth or modify simulated cell motion. + ## Dataset presentation lifecycle Opening a scene file, live session, or recording begins a new dataset. Call `DatasetPresentationState.beginDataset()` and `ColonyViewer.beginDataset()` once, then present its first frame with `setFrame(frame, true)` to fit the camera. Ordinary updates, a reset of the same live model, and recording seeks use `setFrame(frame)` without beginning a dataset. Neither simulation time returning to zero nor a changed signal-grid shape identifies a new dataset. diff --git a/viewer/browser/README.md b/viewer/browser/README.md new file mode 100644 index 0000000..3973ca1 --- /dev/null +++ b/viewer/browser/README.md @@ -0,0 +1,22 @@ +# Capsule rendering verification + +`capsules.mjs` exercises the real viewer renderer through a development-server-only test injection. It records the isolated rod and dense colony at near/far camera distances, 48 frames of prescribed colony motion, stable-ID selection through frame reordering/removal, colors, pointer picking, selection geometry, and a selected zero-length cell. The generated fixtures use arbitrary 3D rod directions; no model or scientific output is modified. + +Start this worktree's Vite server and run the script from the repository root with an existing Playwright installation: + +```console +pnpm --dir viewer install --frozen-lockfile +pnpm --dir viewer dev --host 127.0.0.1 --port 4321 +``` + +In another terminal: + +```console +MICROSIMULATOR_PLAYWRIGHT_MODULE=/absolute/path/to/@playwright/test/index.mjs CAPSULE_MODE=corrected EVIDENCE_DIR=/tmp/capsule-corrected node viewer/browser/capsules.mjs +``` + +`VIEWER_URL` selects another server. For a before/after comparison, run the same script against a separate checkout of `b69193b9` served on another port, using `CAPSULE_MODE=baseline` and a different evidence directory. The script changes only the browser's fetched development module to observe the renderer; it does not patch the checkout. Baseline mode records rather than asserts the corrected GPU-resource and highlight invariants. Both runs use the same camera, lighting, colors, timestep sequence, and fixture parameters. + +The performance fixture renders 512 cells after ten warmup frames, measuring 60 complete frame replacements (transform upload, coloring, rendering, and `gl.finish()` GPU synchronization). It reports median/p95 wall times, browser/WebGL renderer, three draw calls, per-mesh vertices/triangles, and geometry/buffer counts. Timings include CPU work and synchronization and are local regression evidence, not a cross-device benchmark. It tracks actual `createBuffer`/`deleteBuffer` calls to detect instance-buffer leaks that `renderer.info.memory.geometries` alone misses. Corrected mode requires stable buffer counts across replacements and zero tracked buffers after an empty frame and viewer disposal. + +Inspect the PNGs and recorded WebM listed in `results.json`. Compare the tangent joins under identical lighting, distinguishing the reproduced ring seams from silhouette tessellation, pixel aliasing, and the fixture's deliberately changing orientation/position. Do not use test/build success alone to declare the visual artifact resolved. diff --git a/viewer/browser/capsules.mjs b/viewer/browser/capsules.mjs new file mode 100644 index 0000000..67a62a5 --- /dev/null +++ b/viewer/browser/capsules.mjs @@ -0,0 +1,347 @@ +// Vite browser regression: same fixture/camera/light for baseline and corrected meshes. +import assert from "node:assert/strict"; +import { mkdir, writeFile } from "node:fs/promises"; + +const { chromium } = await import( + process.env.MICROSIMULATOR_PLAYWRIGHT_MODULE ?? "@playwright/test" +); +const url = process.env.VIEWER_URL ?? "http://127.0.0.1:4321"; +const mode = process.env.CAPSULE_MODE ?? "corrected"; +const evidence = + process.env.EVIDENCE_DIR ?? `/tmp/microsimulator-capsules-${mode}`; +await mkdir(evidence, { recursive: true }); +const browser = await chromium.launch({ headless: true }); +const context = await browser.newContext({ + viewport: { width: 1440, height: 960 }, + recordVideo: { dir: evidence, size: { width: 1440, height: 960 } }, +}); +const page = await context.newPage(); +const errors = []; +page.on("pageerror", (error) => { + errors.push(error.message); + console.error(error.message); +}); +page.on("console", (message) => { + if (message.type() === "error") console.error(message.text()); +}); +page.on("requestfailed", (request) => + console.error(request.url(), request.failure()), +); +await page.route("**/src/colony-viewer.ts*", async (route) => { + const response = await route.fetch(); + const source = await response.text(); + const marker = "this.onSelection = onSelection;"; + assert.equal(source.split(marker).length, 2); + await route.fulfill({ + response, + body: source.replace(marker, `${marker}\nglobalThis.__testViewer = this;`), + }); +}); + +try { + await page.goto(url); + await page.waitForFunction(() => globalThis.__testViewer !== undefined); + await page.evaluate(() => { + const v = globalThis.__testViewer; + v.renderer.setAnimationLoop(null); + v.controls.enableDamping = false; + v.grid.visible = false; + const gl = v.renderer.getContext(); + const create = gl.createBuffer.bind(gl); + const remove = gl.deleteBuffer.bind(gl); + const buffers = new Set(); + gl.createBuffer = () => { + const b = create(); + buffers.add(b); + return b; + }; + gl.deleteBuffer = (b) => { + buffers.delete(b); + return remove(b); + }; + globalThis.__liveBuffers = buffers; + globalThis.__fixture = (count = 1, time = 0) => { + const columns = count === 1 ? 1 : Math.ceil(Math.sqrt(count * 2)); + const rows = Math.ceil(count / columns); + return { + time, + backend: { + kind: "cpu", + name: "synthetic geometry fixture", + device: "host", + deviceIndex: 0, + native: true, + }, + speciesCount: 0, + signalGrid: null, + constraints: { boxes: [], cylinders: [], planes: [], spheres: [] }, + cells: Array.from({ length: count }, (_, i) => { + const angle = 0.4 + 0.09 * Math.sin(i * 0.7 + time); + const tilt = 0.15 * Math.cos(i * 0.9 + time); + return { + id: String(i + 1), + parentId: null, + slot: i, + position: [ + ((i % columns) - (columns - 1) / 2) * 1.9 + + 0.04 * Math.sin(time + i), + (Math.floor(i / columns) - (rows - 1) / 2) * 1.05, + 0, + ], + direction: [ + Math.cos(angle) * Math.cos(tilt), + Math.sin(angle) * Math.cos(tilt), + Math.sin(tilt), + ], + length: count === 1 ? 3 : 1.05, + radius: count === 1 ? 0.6 : 0.4, + growthRate: 0, + cellType: i % 3, + fixed: false, + species: [], + }; + }), + }; + }; + globalThis.__present = (count, time, distance) => { + const frame = globalThis.__fixture(count, time); + v.setFrame(frame, false); + v.grid.visible = false; + v.setCellColors( + frame.cells.map((cell) => + cell.cellType === 0 + ? [0.65, 0.85, 0.72] + : cell.cellType === 1 + ? [0.85, 0.55, 0.3] + : [0.45, 0.65, 0.9], + ), + ); + if (distance !== undefined) { + v.camera.position.set(0, -distance, distance * 0.65); + v.camera.up.set(0, 0, 1); + v.controls.target.set(0, 0, 0); + v.camera.lookAt(v.controls.target); + v.controls.update(); + } + v.renderer.render(v.scene, v.camera); + }; + document.querySelector("#empty-state").hidden = true; + }); + for (const [name, count, distance] of [ + ["isolated-near", 1, 6], + ["isolated-far", 1, 18], + ["colony-near", 64, 12], + ["colony-far", 64, 32], + ]) { + await page.evaluate( + ([n, d]) => globalThis.__present(n, 0, d), + [count, distance], + ); + await page + .locator("#canvas-host") + .screenshot({ path: `${evidence}/${name}.png` }); + } + for (let frame = 0; frame < 48; frame++) { + await page.evaluate((time) => { + globalThis.__present(64, time, 12); + return new Promise(requestAnimationFrame); + }, frame / 24); + if (frame % 12 === 0) { + await page + .locator("#canvas-host") + .screenshot({ path: `${evidence}/moving-${frame}.png` }); + } + } + const metrics = await page.evaluate(async () => { + const v = globalThis.__testViewer; + const samples = []; + const resources = []; + for (let frame = 0; frame < 70; frame++) { + await new Promise(requestAnimationFrame); + const start = performance.now(); + globalThis.__present(512, frame / 24, 60); + // Complete GPU work so timings compare actual rendering, not only enqueue time. + v.renderer.getContext().finish(); + if (frame >= 10) samples.push(performance.now() - start); + if (frame === 10 || frame === 69) { + resources.push({ + geometries: v.renderer.info.memory.geometries, + buffers: globalThis.__liveBuffers.size, + }); + } + } + const info = { ...v.renderer.info.render }; + const gl = v.renderer.getContext(); + const debug = gl.getExtension("WEBGL_debug_renderer_info"); + const device = debug + ? gl.getParameter(debug.UNMASKED_RENDERER_WEBGL) + : gl.getParameter(gl.RENDERER); + const vertices = v.cellMeshes.map( + (mesh) => mesh.geometry.getAttribute("position").count, + ); + const triangles = v.cellMeshes.map((mesh) => mesh.geometry.index.count / 3); + samples.sort((a, b) => a - b); + v.setFrame(globalThis.__fixture(0), false); + v.renderer.render(v.scene, v.camera); + return { + samples: samples.length, + device, + medianMilliseconds: samples[30], + p95Milliseconds: samples[57], + render: info, + vertices, + triangles, + resources, + emptyBuffers: globalThis.__liveBuffers.size, + }; + }); + + const behavior = await page.evaluate(() => { + const v = globalThis.__testViewer; + globalThis.__present(2, 0, 6); + v.selectCell(1); + const selected = v.selectedCellId; + const frame = globalThis.__fixture(2, 1); + frame.cells.reverse(); + frame.cells = frame.cells.map((cell, slot) => ({ ...cell, slot })); + v.setFrame(frame, false); + v.setCellColors([ + [1, 0, 0], + [0, 1, 0], + ]); + const retained = v.selectedCellId; + const highlights = v.highlight.children.map((mesh) => { + mesh.updateMatrixWorld(true); + return { matrix: mesh.matrixWorld.elements.slice() }; + }); + const colors = v.cellMeshes.map((mesh) => + Array.from(mesh.instanceColor.array), + ); + v.selectCell(null); + const center = v.camera.position + .clone() + .fromArray(frame.cells[0].position) + .project(v.camera); + const canvas = v.renderer.domElement.getBoundingClientRect(); + v.renderer.render(v.scene, v.camera); + return { + selected, + retained, + highlights, + colors, + point: { + x: canvas.x + ((center.x + 1) * canvas.width) / 2, + y: canvas.y + ((1 - center.y) * canvas.height) / 2, + }, + }; + }); + assert.equal(behavior.selected, "2"); + assert.equal(behavior.retained, "2"); + for (const colors of behavior.colors) + assert.deepEqual(colors, [1, 0, 0, 0, 1, 0]); + await page.mouse.click(behavior.point.x, behavior.point.y); + assert.equal( + await page.evaluate(() => globalThis.__testViewer.selectedCellId), + "2", + ); + const highlightError = await page.evaluate(() => { + const v = globalThis.__testViewer; + const selected = v.cells.find((cell) => cell.id === v.selectedCellId); + const center = v.camera.position.clone().fromArray(selected.position); + const axis = center.clone().fromArray(selected.direction).normalize(); + let error = 0; + for (const mesh of v.highlight.children) { + mesh.updateMatrixWorld(true); + const positions = mesh.geometry.getAttribute("position"); + for (let i = 0; i < positions.count; i++) { + const point = center + .clone() + .fromBufferAttribute(positions, i) + .applyMatrix4(mesh.matrixWorld); + const along = point.clone().sub(center).dot(axis); + const nearest = center + .clone() + .addScaledVector( + axis, + Math.max( + -selected.length / 2, + Math.min(selected.length / 2, along), + ), + ); + error = Math.max( + error, + Math.abs(point.distanceTo(nearest) - selected.radius * 1.08), + ); + } + } + v.renderer.render(v.scene, v.camera); + return error; + }); + if (mode === "corrected") assert.ok(highlightError < 1e-6); + await page + .locator("#canvas-host") + .screenshot({ path: `${evidence}/selected.png` }); + const removed = await page.evaluate(() => { + const v = globalThis.__testViewer; + v.setFrame(globalThis.__fixture(1), false); + return { selected: v.selectedCellId, visible: v.highlight.visible }; + }); + assert.deepEqual(removed, { selected: null, visible: false }); + if (mode === "corrected") { + const zero = await page.evaluate(() => { + const v = globalThis.__testViewer; + const frame = globalThis.__fixture(1); + frame.cells[0].length = 0; + v.setFrame(frame, false); + v.setCellColors([[0.6, 0.8, 0.7]]); + v.selectCell(0); + v.renderer.render(v.scene, v.camera); + return v.highlight.children.every((mesh) => + mesh.matrix.elements.every(Number.isFinite), + ); + }); + assert.equal(zero, true, "zero-length selection has finite transforms"); + await page + .locator("#canvas-host") + .screenshot({ path: `${evidence}/zero-length-selected.png` }); + } + if (mode === "corrected") { + assert.equal( + metrics.render.calls, + 3, + "cell rendering stays at three instanced draws", + ); + assert.equal( + metrics.resources[0].buffers, + metrics.resources[1].buffers, + "frame replacement must free old instance buffers", + ); + assert.equal( + metrics.emptyBuffers, + 0, + "empty frame must free cell GPU buffers", + ); + } + assert.deepEqual(errors, []); + const disposedBuffers = await page.evaluate(() => { + globalThis.__testViewer.dispose(); + return globalThis.__liveBuffers.size; + }); + if (mode === "corrected") assert.equal(disposedBuffers, 0); + const report = { + mode, + browser: browser.version(), + fixture: + "512 cells; 10 warmup + 60 measured replacement/color/render/gl.finish frames; fixed camera and lighting; 64-cell prescribed motion video", + metrics, + behavior, + highlightError, + disposedBuffers, + video: await page.video().path(), + }; + await writeFile(`${evidence}/results.json`, JSON.stringify(report, null, 2)); + console.log(JSON.stringify(report, null, 2)); +} finally { + await context.close(); + await browser.close(); +} 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 +