forked from erp-dev/erp
46 lines
1.2 KiB
Python
46 lines
1.2 KiB
Python
from django.http.response import StreamingHttpResponse
|
||
from rest_framework.decorators import api_view, permission_classes
|
||
from rest_framework.response import Response
|
||
from contextlib import suppress
|
||
import asyncio
|
||
|
||
|
||
_connections = set()
|
||
|
||
|
||
@api_view(['POST'])
|
||
@permission_classes([])
|
||
def create_sse_event(request):
|
||
"""
|
||
创建一个简单的SSE事件流响应,用于测试和演示目的。
|
||
"""
|
||
|
||
def sse_stream():
|
||
queue = asyncio.Queue()
|
||
_connections.add(queue)
|
||
|
||
try:
|
||
while True:
|
||
data = queue.get()
|
||
yield f"data: {data}\n\n"
|
||
finally:
|
||
_connections.remove(queue)
|
||
|
||
response = StreamingHttpResponse(sse_stream(), content_type='text/event-stream')
|
||
response['Cache-Control'] = 'no-cache'
|
||
return response
|
||
|
||
|
||
@api_view(['POST'])
|
||
@permission_classes([])
|
||
def push_sse_event(request):
|
||
"""
|
||
向所有连接的客户端广播一个SSE事件。
|
||
请求体应包含一个 'message' 字段,表示要发送的消息内容。
|
||
"""
|
||
message = request.data.get('message', 'Hello, SSE!')
|
||
for q in list(_connections):
|
||
with suppress(asyncio.QueueFull):
|
||
q.put_nowait(message)
|
||
return Response({'status': 'message sent'})
|