mirror of
https://github.com/kennethreitz/responder.git
synced 2026-07-21 18:39:29 +00:00
83 lines
1.8 KiB
Python
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()
|