diff --git a/docs/architecture/04-server-api-endpoints.md b/docs/architecture/04-server-api-endpoints.md index 830b2deb7..574f1385e 100644 --- a/docs/architecture/04-server-api-endpoints.md +++ b/docs/architecture/04-server-api-endpoints.md @@ -142,7 +142,7 @@ Rationale: `server.py` owns process startup/lifespan/static mounting, while `gat ### WebSocket protocol -`/ws` is the live browser delivery channel, not a durable state owner: queue, task, Project, review, skill, settings, cost, and update modules persist their own truth, and REST/history endpoints reconstruct it after reload or disconnection. `gateway/contracts.py` describes the frozen envelope shapes and message-type index for Python/JavaScript parity; `gateway/ws.py` performs the actual transport checks — incoming text must decode to a JSON object, extension types must parse as an owned namespace, and built-in `chat` or `command` frames must carry a non-empty payload before entering the message bridge. A `chat` frame's acceptance (ingress lock, locked canonical row, enqueue, echo) runs off the event loop through `gateway._helpers.run_sync_to_completion`, so a held lock or a slow append never stalls other handlers, while a `command` frame stays an inline queue put. The authentication middleware above governs socket admission; the public socket never receives the Host Service token or the owned Claudexor daemon token. +`/ws` is the live browser delivery channel, not a durable state owner: queue, task, Project, review, skill, settings, cost, and update modules persist their own truth, and REST/history endpoints reconstruct it after reload or disconnection. `gateway/contracts.py` describes the frozen envelope shapes and message-type index for Python/JavaScript parity; `gateway/ws.py` performs the actual transport checks — incoming text must decode to a JSON object, extension types must parse as an owned namespace, and built-in `chat` or `command` frames must carry a non-empty payload before entering the message bridge. A `chat` frame's acceptance (ingress lock, locked canonical row, enqueue, echo) runs off the event loop through `gateway._helpers.run_sync_to_completion`, so a held lock or a slow append never stalls other handlers. Each socket chains its chat acceptances so they settle in receive order, but its receive loop never awaits them: a `command` frame (Panic and Restart share the SPA's one socket with chat) stays an inline queue put, admitted even while that socket's earlier chat is held. Disconnect or cancellation returns only after every received chat settled (`settle_to_completion`). The authentication middleware above governs socket admission; the public socket never receives the Host Service token or the owned Claudexor daemon token. The browser constructs one socket for the whole SPA. Feature modules subscribe before connection, and the initial complete Project chat-id set is fetched before the first open so an early Project frame cannot be mistaken for Main traffic. `ws.on(type, listener)` stores listeners in insertion-ordered sets and returns a disposer; emission uses a listener snapshot, so a listener added during dispatch does not receive the current frame and disposing one listener cannot skip its neighbor. Every decoded frame first reaches the generic `message` event and then its type-specific event, which lets Widgets consume reviewed namespaced events without duplicating the socket. diff --git a/docs/architecture/05-supervisor-loop.md b/docs/architecture/05-supervisor-loop.md index b53235a0e..3d4229e29 100644 --- a/docs/architecture/05-supervisor-loop.md +++ b/docs/architecture/05-supervisor-loop.md @@ -63,6 +63,6 @@ Startup and throttled maintenance reconcile three residue classes, the ~600 s pa Cooperative project checkpointing has two equivalent quiescence triggers: a host-minted genesis or cooperative tree is checked when its root settles with no live descendants, and again when the last child settles beneath an already-terminal root — a root-scope budget stop terminalizes the root before its children, so a root-only trigger would see a live tree once and never return. The bounded git chain runs on a daemon thread, revalidates quiescence under the queue lock immediately before mutation, and replays a trigger that arrives during an in-flight check. Only host-minted project roots are eligible: owner-attached folders are never auto-committed, credential-shaped files stay excluded and disclosed, and every material success, skip or error receives a durable receipt. -Bridge intake is batched: `LocalChatBridge.get_updates` blocks only for the first item, then drains a bounded ready snapshot. Every update keeps a monotonic id and its own reply transport (`activate_update_transport`); malformed items are logged and skipped individually. Web owner ingress durably writes the canonical `chat.jsonl` row with `log_chat(require_write=True, ensure_record_boundary=True)` (as the named ingress does, so a torn prior tail cannot swallow it) before queue handoff or `ingress_accepted` echo, passing that row as the dequeue witness. Each async ingress-lock caller (`gateway/ws.py` chat, Host Service delivery, late quiz answers) runs via `gateway._helpers.run_sync_to_completion` off the ASGI loop: unrelated HTTP and sockets stay responsive, later frames on a socket remain ordered, and cancellation settles row → queue → echo. A command on another socket remains an inline queue put. A failed handler requeues the unprocessed tail with its ids for the next read; accepted `/restart` does likewise, while `/panic` never delays the hard stop for handback. Requeue is memory-only and lost on process exit; accepted rows on disk survive without automatic replay. +Bridge intake is batched: `LocalChatBridge.get_updates` blocks only for the first item, then drains a bounded ready snapshot. Every update keeps a monotonic id and its own reply transport (`activate_update_transport`); malformed items are logged and skipped individually. Web owner ingress durably writes the canonical `chat.jsonl` row with `log_chat(require_write=True, ensure_record_boundary=True)` (as the named ingress does, so a torn prior tail cannot swallow it) before queue handoff or `ingress_accepted` echo, passing that row as the dequeue witness. Each async ingress-lock caller (`gateway/ws.py` chat, Host Service delivery, late quiz answers) runs via `gateway._helpers.run_sync_to_completion` off the ASGI loop: unrelated HTTP and sockets stay responsive, a socket's later chats remain ordered, and cancellation settles row → queue → echo. A command on any socket remains an inline queue put. A failed handler requeues the unprocessed tail with its ids for the next read; accepted `/restart` does likewise, while `/panic` never delays the hard stop for handback. Requeue is memory-only and lost on process exit; accepted rows on disk survive without automatic replay. The bridge recognizes `/panic`, `/restart`, `/review`, `/evolve [on|off]`, `/bg [start|stop|status]` and `/status`; all other text enters ordinary agent routing. External transports may invoke these commands only with positive owner identity and a transport-specific owner-chat binding, and the commands reuse runtime-mode, queue, cancellation and typed-result authority rather than implementing parallel control paths. Runtime logs rotate on the same supervisor tick and archive readers preserve their retained timelines; only explicitly isolated devtool roots may use the narrow rotation sentinel from §1. diff --git a/docs/development/06-rules-by-change-class.md b/docs/development/06-rules-by-change-class.md index 56accaa6e..542dfa3f8 100644 --- a/docs/development/06-rules-by-change-class.md +++ b/docs/development/06-rules-by-change-class.md @@ -106,7 +106,7 @@ Run roots are append-only outside `repo/` and live `data/`; the focused contract than waiting for a growing file to reach EOF. HTTP admission and materialization run their whole blocking operation off the event loop (`gateway._helpers.run_sync_to_completion`); every async ingress-lock - caller uses it for locked row → queue → echo. Cancellation settles before release; other sockets + caller uses it for locked row → queue → echo. Cancellation settles before release; receive loops stay responsive. Cancelling an HTTP waiter never cancels the admitted task. Directory exports carry a complete relative member/size/SHA manifest plus a streamed ZIP (outputs above 50 MiB included); a changed file or missing member is an explicit capture failure, while genesis LISTING is diff --git a/ouroboros/gateway/_helpers.py b/ouroboros/gateway/_helpers.py index d1bbd044d..1b0e8afae 100644 --- a/ouroboros/gateway/_helpers.py +++ b/ouroboros/gateway/_helpers.py @@ -35,25 +35,33 @@ async def run_sync_to_completion(function, /, *args, **kwargs): try: return await asyncio.shield(worker) except asyncio.CancelledError: - # Starlette streams use level cancellation; shield that scope while - # also tolerating repeated raw asyncio Task.cancel() calls. - with anyio.CancelScope(shield=True): - while not worker.done(): - try: - await asyncio.shield(worker) - except asyncio.CancelledError: - pass - except Exception: - break - try: - worker.result() - except Exception: - logging.getLogger(__name__).debug( - "Request worker failed while cancellation settled", exc_info=True, - ) + await settle_to_completion(worker) raise +async def settle_to_completion(task: asyncio.Task) -> None: + """Wait until an owned task is done, through the caller's cancellation. + + Its failure is logged, not raised; a cancelled task still raises. + """ + # Starlette streams use level cancellation; shield that scope while + # also tolerating repeated raw asyncio Task.cancel() calls. + with anyio.CancelScope(shield=True): + while not task.done(): + try: + await asyncio.shield(task) + except asyncio.CancelledError: + pass + except Exception: + break + try: + task.result() + except Exception: + logging.getLogger(__name__).debug( + "Request worker failed while cancellation settled", exc_info=True, + ) + + def read_rotated_jsonl_entries( live: pathlib.Path, archive_dir: pathlib.Path, @@ -180,5 +188,5 @@ def stage_initial_task_attachments( __all__ = ( "coerce_bool", "coerce_int", "iter_jsonl_objects", "json_error", "json_exception", "read_rotated_jsonl_entries", "request_json_or", "request_drive_root", "request_repo_dir", - "stage_initial_task_attachments", "run_sync_to_completion", + "stage_initial_task_attachments", "run_sync_to_completion", "settle_to_completion", ) diff --git a/ouroboros/gateway/ws.py b/ouroboros/gateway/ws.py index 82a83e810..a3678b80c 100644 --- a/ouroboros/gateway/ws.py +++ b/ouroboros/gateway/ws.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio +import functools import inspect import json import logging @@ -13,7 +14,7 @@ from typing import Any from starlette.websockets import WebSocket, WebSocketDisconnect from ouroboros.config import DATA_DIR -from ouroboros.gateway._helpers import run_sync_to_completion +from ouroboros.gateway._helpers import run_sync_to_completion, settle_to_completion from ouroboros.utils import utc_now_iso log = logging.getLogger(__name__) @@ -315,12 +316,45 @@ async def _dispatch_extension_message( return True +def _initialization_notice() -> str: + return json.dumps({ + "type": "chat", + "role": "system", "system_type": "initialization_notice", + "content": "⚠️ System is still initializing. Please wait a moment and try again.", + "ts": utc_now_iso(), + }) + + +async def _accept_chat_after(websocket: WebSocket, previous: asyncio.Task | None, accept) -> None: + """Run one chat frame's acceptance once this socket's previous one settled. + + The web acceptance takes the single host's ingress lock and a locked durable + append (log_chat(require_write=True)); either may wait behind a skill delivery + scanning retained chat or a slow disk, so it runs off the ASGI loop, and the + receive loop never awaits it: a command frame on the same socket (Panic, + Restart) is admitted meanwhile. ``run_sync_to_completion`` keeps custody of + row → queue → echo; a failure answers the sender with the initialization notice. + """ + if previous is not None: + await asyncio.wait((previous,)) # its failure was its own notice + try: + await run_sync_to_completion(accept) + except Exception: + try: + await websocket.send_text(_initialization_notice()) + except Exception: + log.debug("WebSocket closed before its chat failure notice", exc_info=True) + + async def ws_endpoint(websocket: WebSocket) -> None: await websocket.accept() with _ws_lock: _ws_clients.append(websocket) total = len(_ws_clients) log.info("WebSocket client connected (total: %d)", total) + # Tail of this socket's chat acceptances: each waits for the one before it, + # so its chat frames settle in receive order while the loop keeps receiving. + accepting: asyncio.Task | None = None try: while True: data = await websocket.receive_text() @@ -391,15 +425,7 @@ async def ws_endpoint(websocket: WebSocket) -> None: # land in chat.jsonl at the canonical-row writer. client_surface["received_at"] = utc_now_iso() task_metadata["client_surface"] = client_surface - # The web acceptance takes the single host's ingress lock and - # a locked durable append (log_chat(require_write=True)); either - # may wait behind a skill delivery scanning retained chat or a - # slow disk, so it must not run on the ASGI loop. The settled - # worker wait keeps its custody: a cancelled socket task still - # lets the canonical row, the queue item and the echo complete - # in that order, and this socket's frames stay ordered because - # the next receive waits for it. - await run_sync_to_completion( + accept = functools.partial( bridge.ui_send, payload, sender_session_id=str(msg.get("sender_session_id", "") or ""), @@ -411,17 +437,14 @@ async def ws_endpoint(websocket: WebSocket) -> None: chat_id=thread_id, project_id=str(msg.get("project_id", "") or ""), ) + accepting = asyncio.create_task(_accept_chat_after(websocket, accepting, accept)) else: - # A command is a queue put only (no lock, no durable write): - # it stays inline once received; another socket is not blocked by owner ingress. + # A command is a queue put only (no lock, no durable write). It is + # admitted on receipt, never behind this socket's pending chat + # acceptances: Panic and Restart ride it on the SPA's one socket. bridge.ui_send(payload, broadcast=False) except Exception: - await websocket.send_text(json.dumps({ - "type": "chat", - "role": "system", "system_type": "initialization_notice", - "content": "⚠️ System is still initializing. Please wait a moment and try again.", - "ts": utc_now_iso(), - })) + await websocket.send_text(_initialization_notice()) except WebSocketDisconnect: pass except Exception as exc: @@ -434,6 +457,10 @@ async def ws_endpoint(websocket: WebSocket) -> None: pass total = len(_ws_clients) log.info("WebSocket client disconnected (total: %d)", total) + if accepting is not None: + # Received chat frames keep custody through disconnect or cancellation: + # the socket task returns only after each settled row → queue → echo. + await settle_to_completion(accepting) __all__ = [ diff --git a/tests/test_ws_ingress_off_loop.py b/tests/test_ws_ingress_off_loop.py index 98a9e78b4..81d6e8485 100644 --- a/tests/test_ws_ingress_off_loop.py +++ b/tests/test_ws_ingress_off_loop.py @@ -8,8 +8,10 @@ event loop with it: every other HTTP and WebSocket handler waited, Stop and Panic included. The acceptance now runs off the loop through ``gateway._helpers.run_sync_to_completion`` and keeps its custody: the row, the queue item and the echo still complete in that order, even when the socket task -is cancelled while the acceptance stands on the lock. A command on another -socket stays an inline queue put; a later frame on this socket waits its turn. +is cancelled while the acceptance stands on the lock. A later chat frame on +this socket waits its turn, but the receive loop never waits: the SPA sends +Panic and Restart on the SAME socket as chat, so a command frame is admitted +inline while that socket's chat acceptance is still held (review F2). A late quiz answer (``POST /api/decisions`` on an expired card) is forwarded through the named ingress, which waits for that same lock — behind a socket @@ -24,10 +26,12 @@ import json import threading import time +import anyio import pytest from starlette.applications import Starlette from starlette.routing import Route, WebSocketRoute from starlette.testclient import TestClient +from starlette.websockets import WebSocketDisconnect from ouroboros.gateway.state import api_health from ouroboros.gateway.ws import ws_endpoint @@ -216,6 +220,272 @@ def test_an_acceptance_failure_still_answers_the_socket(bridge): assert echoes == [] and bridge.get_updates(offset=0, timeout=0) == [] + +def _record_handoffs(bridge, release: threading.Event | None = None): + """Instrument the bridge callback the socket hands every frame to: record each + chat acceptance entering it (held until ``release`` when given) and each command + admitted. Nothing consumes the queue, so an admitted /panic is never executed.""" + handoffs: list[tuple[str, str]] = [] + chat_entered = threading.Event() + original = bridge.ui_send + + def ui_send(text, **kwargs): + if kwargs.get("broadcast") is False: + handoffs.append(("command", text)) + else: + handoffs.append(("chat", kwargs["client_message_id"])) + chat_entered.set() + if release is not None: + assert release.wait(5), "the test never released the held acceptance" + return original(text, **kwargs) + + bridge.ui_send = ui_send + return handoffs, chat_entered + + +def _wait_for(predicate, timeout: float = 3.0) -> bool: + deadline = time.monotonic() + timeout + while not predicate() and time.monotonic() < deadline: + time.sleep(0.01) + return predicate() + + +def _command_frame(cmd: str) -> str: + return json.dumps({"type": "command", "cmd": cmd}) + + +def _queued(bridge) -> list[str]: + with bridge._inbox.mutex: + return [item["client_message_id"] or item["text"] for item in bridge._inbox.queue] + + +def test_a_same_socket_panic_is_admitted_while_its_chat_acceptance_is_held(bridge): + """Consumer regression (TZ-1 PR1 review F2): ``confirmAndSendPanic`` uses the one + SPA socket that also carries chat. With that socket's chat acceptance standing on + a held ingress lock, its next frame — Panic — must still be received and admitted + on one event loop; the Emergency Stop never waits for ordinary acceptance work.""" + held, release, timed_out = threading.Event(), threading.Event(), threading.Event() + frames, echoed = _witness_custody(bridge, chat_echoes=1) + handoffs, chat_entered = _record_handoffs(bridge) + + def hold_ingress(): # a skill delivery mid-scan under the single host's ingress lock + with message_bus._INGRESS_LOCK: + held.set() + if not release.wait(timeout=6.0): + timed_out.set() + + holder = threading.Thread(target=hold_ingress, name="ingress-holder", daemon=True) + holder.start() + assert held.wait(5) + app = Starlette(routes=[WebSocketRoute("/ws", ws_endpoint)]) + try: + with TestClient(app) as client, client.websocket_connect("/ws") as owner: + owner.send_text(_chat_frame("held-1")) + assert chat_entered.wait(5), "the chat frame never reached the bridge" + owner.send_text(_command_frame("/panic")) + assert _wait_for(lambda: ("command", "/panic") in handoffs), ( + "Panic waited behind this socket's held chat acceptance") + assert not timed_out.is_set() and message_bus._INGRESS_LOCK.locked() + assert frames == [] and _row_ids(bridge.drive) == [] and _queued(bridge) == ["/panic"] + release.set() + holder.join(5) + assert echoed.wait(5), frames + finally: + release.set() + assert handoffs == [("chat", "held-1"), ("command", "/panic")] + [(echo, queued, rows)] = frames + assert echo["client_message_id"] == "held-1" and echo["ingress_accepted"] is True + assert queued == ["/panic", "held-1"] and rows == ["held-1"] # row → queue → echo + + +def test_same_socket_chat_frames_settle_in_receive_order_while_commands_pass(bridge): + """Later chat frames on the socket wait for the held one — they never race it for + the lock — while commands between them are admitted on receipt. After release the + rows, queue items and echoes follow receive order, each echo after its row and item.""" + release = threading.Event() + frames, echoed = _witness_custody(bridge, chat_echoes=3) + handoffs, chat_entered = _record_handoffs(bridge, release) + app = Starlette(routes=[WebSocketRoute("/ws", ws_endpoint)]) + try: + with TestClient(app) as client, client.websocket_connect("/ws") as owner: + owner.send_text(_chat_frame("c-1", "first")) + assert chat_entered.wait(5) + owner.send_text(_chat_frame("c-2", "second")) + owner.send_text(_command_frame("/restart")) + owner.send_text(_chat_frame("c-3", "third")) + owner.send_text(_command_frame("/status")) # its admission proves c-3 was received + assert _wait_for(lambda: len(handoffs) == 3), handoffs + assert handoffs == [("chat", "c-1"), ("command", "/restart"), ("command", "/status")] + assert frames == [] and _row_ids(bridge.drive) == [] + release.set() + assert echoed.wait(5), frames + finally: + release.set() + assert handoffs[3:] == [("chat", "c-2"), ("chat", "c-3")] + assert _row_ids(bridge.drive) == ["c-1", "c-2", "c-3"] + assert [frame["client_message_id"] for frame, *_ in frames] == ["c-1", "c-2", "c-3"] + for frame, queued, rows in frames: + assert frame["client_message_id"] in queued and frame["client_message_id"] in rows, frames + updates = bridge.get_updates(offset=0, timeout=1) + assert [u["message"]["text"] for u in updates] == ["/restart", "/status", "first", "second", "third"] + + +def test_a_disconnect_settles_every_received_chat_before_the_socket_task_returns(bridge): + """The owner closes the socket (TestClient then cancels the session's anyio scope) + while its first chat is held and a second waits behind it: the socket task returns + only after both settled row → queue → echo, in receive order.""" + release = threading.Event() + frames, echoed = _witness_custody(bridge, chat_echoes=2) + handoffs, chat_entered = _record_handoffs(bridge, release) + app = Starlette(routes=[WebSocketRoute("/ws", ws_endpoint)]) + with TestClient(app) as client: + owner = client.websocket_connect("/ws") + owner.__enter__() + closer = threading.Thread(target=owner.__exit__, args=(None, None, None), daemon=True) + try: + owner.send_text(_chat_frame("d-1", "first")) + assert chat_entered.wait(5) + owner.send_text(_chat_frame("d-2", "second")) + owner.send_text(_command_frame("/status")) # its admission proves d-2 was received + assert _wait_for(lambda: ("command", "/status") in handoffs), handoffs + closer.start() # disconnect, then the session's anyio scope is cancelled + closer.join(0.3) + assert closer.is_alive(), "the socket task returned with received chat frames unsettled" + assert frames == [] and _row_ids(bridge.drive) == [] + release.set() + closer.join(5) + assert not closer.is_alive() + assert echoed.wait(5), frames + finally: # a failed assertion still closes the session instead of hanging the client + release.set() + if closer.ident is None: + closer.start() + closer.join(10) + assert _row_ids(bridge.drive) == ["d-1", "d-2"] + assert [frame["client_message_id"] for frame, *_ in frames] == ["d-1", "d-2"] + for frame, queued, rows in frames: + assert frame["client_message_id"] in queued and frame["client_message_id"] in rows, frames + + +class _ClosingSocket(_OpenSocket): + """A socket whose client disconnects after its frames; later sends fail.""" + + def __init__(self, frames, *, disconnect: bool): + super().__init__(frames) + self.disconnect = disconnect + self.closed = False + + async def receive_text(self): + if self.frames or not self.disconnect: + return await super().receive_text() + self.closed = True + raise WebSocketDisconnect(1001) + + async def send_text(self, text): + if self.closed: + raise RuntimeError("websocket is closed") + await super().send_text(text) + + +@pytest.mark.parametrize("disconnect", [False, True]) +def test_a_failed_acceptance_answers_in_order_and_the_next_chat_still_lands(bridge, disconnect): + """A refused acceptance answers its sender with the initialization notice before + the next chat on the socket is accepted. If the socket is already gone, the + undeliverable notice is dropped quietly, the next chat still lands, and the + socket task ends cleanly — no task exception is left unretrieved.""" + release = threading.Event() + timeline: list[str] = [] + bridge._broadcast_fn = lambda payload: timeline.append(f"echo:{payload['client_message_id']}") + original = bridge.ui_send + + def refuse_first(text, **kwargs): + if kwargs["client_message_id"] == "refused-1": + assert release.wait(5) + raise OSError("disk refused") + return original(text, **kwargs) + + bridge.ui_send = refuse_first + socket = _ClosingSocket([_chat_frame("refused-1"), _chat_frame("ok-2", "next")], disconnect=disconnect) + sent = socket.send_text + + async def send_text(text): + await sent(text) + timeline.append(f"sent:{json.loads(text)['system_type']}") + + socket.send_text = send_text + unretrieved: list[dict] = [] + + async def main(): + asyncio.get_running_loop().set_exception_handler(lambda _loop, context: unretrieved.append(context)) + task = asyncio.create_task(ws_endpoint(socket)) + await asyncio.sleep(0.05) + assert not task.done() and timeline == [] + release.set() + deadline = time.monotonic() + 5.0 + while "echo:ok-2" not in timeline and time.monotonic() < deadline: + await asyncio.sleep(0.01) + if disconnect: + await asyncio.wait_for(task, 5) # returns cleanly, nothing re-raised + else: + task.cancel() + with pytest.raises(asyncio.CancelledError): + await asyncio.wait_for(task, 5) + + asyncio.run(main()) + expected = ["echo:ok-2"] if disconnect else ["sent:initialization_notice", "echo:ok-2"] + assert timeline == expected + assert _row_ids(bridge.drive) == ["ok-2"] and not unretrieved + [update] = bridge.get_updates(offset=0, timeout=1) + assert update["message"]["text"] == "next" + + +@pytest.mark.parametrize("mode", ["asyncio", "anyio"]) +def test_a_cancelled_socket_settles_every_received_chat_in_order(bridge, mode): + """Cancel the socket task — raw ``Task.cancel()`` twice, or anyio's level + cancellation — while its first chat is held and a second waits behind it: the + task stays in custody until both settled in receive order.""" + release = threading.Event() + frames, _echoed = _witness_custody(bridge, chat_echoes=2) + handoffs, chat_entered = _record_handoffs(bridge, release) + socket = _OpenSocket([_chat_frame("x-1", "first"), _chat_frame("x-2", "second")]) + scopes: list[anyio.CancelScope] = [] + + async def endpoint(): + with anyio.CancelScope() as scope: + scopes.append(scope) + await ws_endpoint(socket) + + async def main(): + task = asyncio.create_task(endpoint()) + assert await asyncio.to_thread(chat_entered.wait, 5) + await asyncio.sleep(0.01) # the socket is now waiting for a third frame + if mode == "anyio": + scopes[0].cancel() + else: + task.cancel() + await asyncio.sleep(0) + task.cancel() + await asyncio.sleep(0.05) + assert not task.done(), "the cancelled socket task abandoned its received chats" + assert frames == [] and handoffs == [("chat", "x-1")] + release.set() + if mode == "anyio": + await asyncio.wait_for(task, 5) # the scope absorbs its own cancellation + else: + with pytest.raises(asyncio.CancelledError): + await asyncio.wait_for(task, 5) + + try: + asyncio.run(main()) + finally: + release.set() + assert handoffs == [("chat", "x-1"), ("chat", "x-2")] + assert _row_ids(bridge.drive) == ["x-1", "x-2"] + assert [frame["client_message_id"] for frame, *_ in frames] == ["x-1", "x-2"] + for frame, queued, rows in frames: + assert frame["client_message_id"] in queued and frame["client_message_id"] in rows, frames + + _LATE_ID = "quiz_late_answer:task-1:q1" _LATE_BODY = {"request_id": "r1", "decision_id": "quiz:task-1:q1", "option_index": 1}