|
1 | 1 | import asyncio |
2 | 2 | import logging |
| 3 | +from datetime import timedelta |
3 | 4 | from typing import ( |
4 | 5 | Any, |
5 | 6 | AsyncIterator, |
|
10 | 11 |
|
11 | 12 | from replit_river.messages import parse_transport_msg |
12 | 13 | from replit_river.rpc import ( |
| 14 | + STREAM_OPEN_BIT, |
13 | 15 | ControlMessageHandshakeRequest, |
14 | 16 | ControlMessageHandshakeResponse, |
15 | 17 | HandShakeStatus, |
16 | 18 | TransportMessage, |
17 | 19 | ) |
18 | 20 | from replit_river.transport_options import TransportOptions |
19 | 21 | from replit_river.v2.client import Client |
20 | | -from replit_river.v2.session import STREAM_CANCEL_BIT |
| 22 | +from replit_river.v2.session import STREAM_CANCEL_BIT, STREAM_CLOSED_BIT |
21 | 23 | from tests.v2.fixtures.raw_ws_server import OuterPayload, WsServerFixture |
22 | 24 |
|
23 | 25 |
|
| 26 | +async def test_rpc_cancel(ws_server: WsServerFixture) -> None: |
| 27 | + (urimeta, recv, conn) = ws_server |
| 28 | + |
| 29 | + client = Client( |
| 30 | + client_id="CLIENT1", |
| 31 | + server_id="SERVER", |
| 32 | + transport_options=TransportOptions(), |
| 33 | + uri_and_metadata_factory=urimeta, |
| 34 | + ) |
| 35 | + |
| 36 | + connecting = asyncio.create_task(client.ensure_connected()) |
| 37 | + request_msg = parse_transport_msg(await recv.get()) |
| 38 | + |
| 39 | + assert not isinstance(request_msg, str) |
| 40 | + assert (serverconn := conn()) |
| 41 | + handshake_request: ControlMessageHandshakeRequest[None] = ( |
| 42 | + ControlMessageHandshakeRequest(**request_msg.payload) |
| 43 | + ) |
| 44 | + |
| 45 | + handshake_resp = ControlMessageHandshakeResponse( |
| 46 | + status=HandShakeStatus( |
| 47 | + ok=True, |
| 48 | + ), |
| 49 | + ) |
| 50 | + handshake_request.sessionId |
| 51 | + |
| 52 | + msg = TransportMessage( |
| 53 | + from_=request_msg.from_, |
| 54 | + to=request_msg.to, |
| 55 | + streamId=request_msg.streamId, |
| 56 | + controlFlags=0, |
| 57 | + id=nanoid.generate(), |
| 58 | + seq=0, |
| 59 | + ack=0, |
| 60 | + payload=handshake_resp.model_dump(), |
| 61 | + ) |
| 62 | + packed = msgpack.packb( |
| 63 | + msg.model_dump(by_alias=True, exclude_none=True), datetime=True |
| 64 | + ) |
| 65 | + await serverconn.send(packed) |
| 66 | + |
| 67 | + sent_waiter = asyncio.Event() |
| 68 | + |
| 69 | + async def handle_server_messages() -> None: |
| 70 | + request_msg = parse_transport_msg(await recv.get()) |
| 71 | + assert not isinstance(request_msg, str) |
| 72 | + |
| 73 | + logging.debug("request_msg: %r", repr(request_msg)) |
| 74 | + |
| 75 | + assert request_msg.payload.get("payload", {}).get("hello") == "world" |
| 76 | + logging.debug("Found a hello:world %r", repr(request_msg)) |
| 77 | + |
| 78 | + sent_waiter.set() |
| 79 | + |
| 80 | + assert request_msg.controlFlags == STREAM_OPEN_BIT | STREAM_CLOSED_BIT |
| 81 | + |
| 82 | + cancel_msg = parse_transport_msg(await recv.get()) |
| 83 | + assert not isinstance(cancel_msg, str) |
| 84 | + assert cancel_msg.controlFlags == STREAM_CANCEL_BIT |
| 85 | + |
| 86 | + server_handler = asyncio.create_task(handle_server_messages()) |
| 87 | + |
| 88 | + rpc_task = asyncio.create_task( |
| 89 | + client.send_rpc( |
| 90 | + "test", |
| 91 | + "bigstream", |
| 92 | + {"ok": True, "payload": {"hello": "world"}}, |
| 93 | + lambda x: x, |
| 94 | + lambda x: x, |
| 95 | + lambda x: x, |
| 96 | + timedelta(seconds=2), |
| 97 | + ) |
| 98 | + ) |
| 99 | + |
| 100 | + # Wait until we've seen at least a few messages from the upload Task |
| 101 | + await sent_waiter.wait() |
| 102 | + |
| 103 | + rpc_task.cancel() |
| 104 | + |
| 105 | + try: |
| 106 | + await rpc_task |
| 107 | + except asyncio.CancelledError: |
| 108 | + pass |
| 109 | + |
| 110 | + await client.close() |
| 111 | + await connecting |
| 112 | + |
| 113 | + # Ensure we're listening to close messages as well |
| 114 | + server_handler.cancel() |
| 115 | + await server_handler |
| 116 | + |
| 117 | + |
24 | 118 | async def test_upload_cancel(ws_server: WsServerFixture) -> None: |
25 | 119 | (urimeta, recv, conn) = ws_server |
26 | 120 |
|
|
0 commit comments