Files
responder/examples/sse_stream.py
T

83 lines
1.8 KiB
Python

"""Server-Sent Events streaming example.
Run it:
responder run examples/sse_stream.py
Try it with:
curl -N http://127.0.0.1:5042/stream
"""
from __future__ import annotations
import asyncio
from collections.abc import AsyncIterator
from pydantic import BaseModel
import responder
class Tick(BaseModel):
number: int
message: str
def create_api(*, event_count: int = 20, delay: float = 0.5) -> responder.API:
api = responder.API(
title="SSE Stream",
version="1.0",
openapi="3.1.0",
docs_route="/docs",
sessions=False,
)
@api.get("/", include_in_schema=False)
def index(req, resp):
resp.html = """
<!DOCTYPE html>
<html>
<body>
<h1>SSE Stream</h1>
<div id="events"></div>
<script>
const source = new EventSource("/stream");
const events = document.getElementById("events");
source.addEventListener("tick", (event) => {
const tick = JSON.parse(event.data);
const p = document.createElement("p");
p.textContent = `${tick.number}. ${tick.message}`;
events.appendChild(p);
});
</script>
</body>
</html>
"""
@api.sse(
"/stream",
heartbeat=15,
operation_id="stream_events",
tags=["events"],
summary="Stream events",
)
async def stream(req, resp) -> AsyncIterator[responder.SSE[Tick]]:
for event_id in range(1, event_count + 1):
yield responder.SSE(
Tick(number=event_id, message=f"Event #{event_id}"),
id=str(event_id),
event="tick",
)
if delay:
await asyncio.sleep(delay)
return api
api = create_api()
if __name__ == "__main__":
api.run()