Skip to content
File

Blob: src/workerd/server/tests/python/asgi-sse/worker.py

python103 lines
1class 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 
55from workers import WorkerEntrypoint
56 
57 
58class 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 
68app = SSEServer()
69 
70 
71async 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