File
Blob: src/workerd/server/tests/python/asgi-sse/worker.py
| 1 | class SSEServer: |
| 2 | def __init__(self): |
| 3 | pass |
| 4 | |
| 5 | async def __call__(self, scope, receive, send) -> None: |
| 6 | scope["app"] = self |
| 7 | |
| 8 | assert scope["type"] in ("http", "websocket", "lifespan") |
| 9 | |
| 10 | if scope["type"] == "lifespan": |
| 11 | message = await receive() |
| 12 | if message["type"] == "lifespan.startup": |
| 13 | await send({"type": "lifespan.startup.complete"}) |
| 14 | return |
| 15 | |
| 16 | elif scope["type"] == "http": |
| 17 | # Receive the request |
| 18 | await receive() |
| 19 | |
| 20 | # Send SSE response headers |
| 21 | await send( |
| 22 | { |
| 23 | "type": "http.response.start", |
| 24 | "status": 200, |
| 25 | "headers": [ |
| 26 | (b"cache-control", b"no-store"), |
| 27 | (b"connection", b"keep-alive"), |
| 28 | (b"content-type", b"text/event-stream; charset=utf-8"), |
| 29 | (b"x-accel-buffering", b"no"), |
| 30 | ], |
| 31 | } |
| 32 | ) |
| 33 | |
| 34 | # Send initial event |
| 35 | await send( |
| 36 | { |
| 37 | "type": "http.response.body", |
| 38 | "body": b"event: endpoint\r\ndata: /messages/?session_id=test123\r\n\r\n", |
| 39 | "more_body": True, |
| 40 | } |
| 41 | ) |
| 42 | |
| 43 | # Send three ping events |
| 44 | for i in range(3): |
| 45 | # In a real app we would wait between events, but in the test we'll send them quickly |
| 46 | await send( |
| 47 | { |
| 48 | "type": "http.response.body", |
| 49 | "body": f": ping - message {i + 1}\r\n\r\n".encode(), |
| 50 | "more_body": i < 2, # last message has more_body=False |
| 51 | } |
| 52 | ) |
| 53 | |
| 54 | |
| 55 | from workers import WorkerEntrypoint |
| 56 | |
| 57 | |
| 58 | class Default(WorkerEntrypoint): |
| 59 | async def fetch(self, request): |
| 60 | import asgi |
| 61 | |
| 62 | return await asgi.fetch(app, request, self.env, self.ctx) |
| 63 | |
| 64 | async def test(self, ctrl): |
| 65 | await test_sse(self.env) |
| 66 | |
| 67 | |
| 68 | app = SSEServer() |
| 69 | |
| 70 | |
| 71 | async def test_sse(env): |
| 72 | # Make a request to our SSE endpoint |
| 73 | response = await env.SELF.fetch("http://example.com/sse") |
| 74 | |
| 75 | # Verify the response has the correct headers for SSE |
| 76 | assert response.headers["content-type"] == "text/event-stream; charset=utf-8" |
| 77 | assert response.headers["cache-control"] == "no-store" |
| 78 | |
| 79 | # Use a simple method to convert the stream to text |
| 80 | from js import TextDecoder |
| 81 | |
| 82 | # Read the stream in a simpler way |
| 83 | reader = response.body.getReader() |
| 84 | content = "" |
| 85 | decoder = TextDecoder.new() |
| 86 | |
| 87 | while True: |
| 88 | result = await reader.read() |
| 89 | if result.done: |
| 90 | break |
| 91 | # Use TextDecoder to convert the chunk to text |
| 92 | chunk_text = decoder.decode(result.value, {"stream": True}) |
| 93 | content += chunk_text |
| 94 | |
| 95 | # Final flush |
| 96 | content += decoder.decode() |
| 97 | # Verify the expected events are in the response |
| 98 | assert "event: endpoint" in content |
| 99 | assert "data: /messages/?session_id=test123" in content |
| 100 | assert ": ping - message 1" in content |
| 101 | assert ": ping - message 2" in content |
| 102 | assert ": ping - message 3" in content |