"""The Docker event stream behind the UI's refresh (F16). Every page used to poll its own endpoint every few seconds for state that only changes when Docker does something — which the daemon already announces. The endpoint existed but nothing consumed it, and it was not usable as it stood: * it forwarded the whole firehose, including three ``exec_*`` events per web terminal session, so a client would have invalidated its cache more often than the polling it replaces; * it never told the client *what* changed, so there was nothing to decide which queries to drop; * and it leaked its reader thread — cancelling the executor future does not interrupt a thread already inside a blocking read, so every page load left one behind holding a socket open. These tests pin the shape of the payload, the filtering, and the teardown. """ from __future__ import annotations import pytest @pytest.fixture def ws_module(): from routers import ws return ws # --------------------------------------------------------------------------- # # Filtering # --------------------------------------------------------------------------- # def test_only_ui_relevant_resource_types_are_requested(ws_module): """Filtered daemon-side, so the bulk never crosses the socket.""" assert set(ws_module._EVENT_TYPES) == {"container", "image", "network", "volume"} @pytest.mark.parametrize( "action", ["exec_create", "exec_start", "exec_die", "attach", "top", "resize"] ) def test_noise_actions_are_ignored(ws_module, action): assert action in ws_module._IGNORED_ACTIONS @pytest.mark.parametrize( "action", ["create", "start", "stop", "die", "destroy", "restart", "health_status", "pull"], ) def test_state_changing_actions_are_not_ignored(ws_module, action): assert action not in ws_module._IGNORED_ACTIONS def test_parameterised_actions_still_match_the_ignore_list(ws_module): """Docker reports these as ``exec_create: /bin/sh``, not a bare verb. Matching the whole string would let every one of them through. """ raw = "exec_create: /bin/sh -c 'ls'" assert raw.split(":")[0] in ws_module._IGNORED_ACTIONS # --------------------------------------------------------------------------- # # The stream, driven against a fake daemon # --------------------------------------------------------------------------- # class _FakeStream: """Stands in for docker-py's CancellableStream. ``close()`` is what actually unblocks the reader thread; this records that it was called, which is the whole point of the teardown test. """ def __init__(self, events): self._events = list(events) self.closed = False def __iter__(self): yield from self._events def close(self): self.closed = True @pytest.fixture def fake_docker(monkeypatch): """Point the endpoint at a canned event stream instead of a daemon.""" import docker_client holder = {} def make(events): stream = _FakeStream(events) holder["stream"] = stream class Client: def events(self, decode=True, filters=None): holder["filters"] = filters return stream monkeypatch.setattr(docker_client, "get_client", lambda: Client()) return holder return make def _event(resource, action, name=None, project=None): return { "Type": resource, "Action": action, "Actor": {"Attributes": {"name": name, "com.docker.compose.project": project}}, } def _collect(client, token, expected): """Read frames off the socket until ``expected`` events have arrived.""" import json frames = [] with client.websocket_connect(f"/ws/events?token={token}") as socket: ready = json.loads(socket.receive_text()) assert ready == {"type": "ready"}, "clients rely on this to know it is live" for _ in range(expected): frames.append(json.loads(socket.receive_text())) return frames def test_an_event_carries_what_the_client_needs(client, admin_token, fake_docker): fake_docker([_event("container", "start", "jellyfin", "jellyfin")]) (frame,) = _collect(client, admin_token, 1) assert frame == { "type": "event", "resource": "container", "action": "start", "container": "jellyfin", "stack": "jellyfin", } def test_noise_is_dropped_before_it_reaches_the_client(client, admin_token, fake_docker): fake_docker( [ _event("container", "exec_create: /bin/sh", "jellyfin"), _event("container", "exec_start: /bin/sh", "jellyfin"), _event("container", "exec_die", "jellyfin"), _event("container", "die", "jellyfin", "jellyfin"), ] ) (frame,) = _collect(client, admin_token, 1) assert frame["action"] == "die", "only the state change should survive" def test_the_daemon_is_asked_to_filter(client, admin_token, fake_docker): holder = fake_docker([_event("container", "start", "x")]) _collect(client, admin_token, 1) assert holder["filters"] == {"type": ["container", "image", "network", "volume"]} def test_the_stream_is_closed_on_disconnect(client, admin_token, fake_docker): """Otherwise the reader thread stays blocked for the life of the process. Cancelling the executor future does not interrupt a thread already inside a blocking read — closing the underlying stream is what does. """ holder = fake_docker([_event("container", "start", "x")]) _collect(client, admin_token, 1) assert holder["stream"].closed, "a disconnected client must not leak its reader" def test_an_ending_stream_is_reported_not_silently_dropped( client, admin_token, fake_docker ): """A daemon restart ends the iterator; the client should hear about it so it reconnects rather than sitting on a dead socket believing it is live.""" import json fake_docker([]) # iterator finishes immediately with client.websocket_connect(f"/ws/events?token={admin_token}") as socket: assert json.loads(socket.receive_text())["type"] == "ready" frame = json.loads(socket.receive_text()) assert frame["type"] == "error" def test_events_require_a_token(client): from starlette.websockets import WebSocketDisconnect with pytest.raises(WebSocketDisconnect) as excinfo: with client.websocket_connect("/ws/events") as socket: socket.receive_text() assert excinfo.value.code == 4401 def test_a_revoked_session_is_refused(client, user_token, db): """The socket resolves against the live user, like every other endpoint.""" from sqlmodel import Session, select from starlette.websockets import WebSocketDisconnect import auth as auth_mod from database import engine from models.user import User with Session(engine) as session: user = session.exec(select(User).where(User.username == "test-user")).one() auth_mod.bump_token_version(user) session.add(user) session.commit() try: with pytest.raises(WebSocketDisconnect) as excinfo: with client.websocket_connect(f"/ws/events?token={user_token}") as socket: socket.receive_text() assert excinfo.value.code == 4401 finally: # Leave the shared fixture user usable for the rest of the session. with Session(engine) as session: user = session.exec(select(User).where(User.username == "test-user")).one() user.token_version = 1 session.add(user) session.commit()