106 lines
3.7 KiB
Python
106 lines
3.7 KiB
Python
"""Real-time endpoints — SSE stream & collaboration WebSocket (ROADMAP #85, tranche 9).
|
|
|
|
Handlers déplacés depuis :mod:`backend.main` sans changement de
|
|
comportement : mêmes chemins (``/api/events``,
|
|
``/ws/collab/{vault}/{path}``), même authentification (Depend pour le SSE,
|
|
manuelle pour le WebSocket — les ``Depends`` FastAPI ne s'exécutent pas sur
|
|
les routes WebSocket).
|
|
|
|
Pas de tags déclarés : assignation par chemin via
|
|
``openapi_docs.tag_for_path`` comme avant (``/api/events`` → System).
|
|
"""
|
|
|
|
import asyncio
|
|
import json as _json
|
|
|
|
from fastapi import APIRouter, Depends, WebSocket
|
|
from fastapi.responses import StreamingResponse
|
|
|
|
from backend.auth.middleware import check_vault_access, require_auth
|
|
from backend.collab import authenticate_websocket, collab_manager
|
|
from backend.services.paths import resolve_safe_path
|
|
from backend.services.vaults import get_vault_root
|
|
from backend.sse import sse_manager
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
@router.get(
|
|
"/api/events",
|
|
response_class=StreamingResponse,
|
|
responses={200: {"content": {"text/event-stream": {}}, "description": "Server-Sent Events stream"}},
|
|
)
|
|
async def api_events(current_user=Depends(require_auth)):
|
|
"""SSE stream for real-time index update notifications.
|
|
|
|
Sends keepalive comments every 30s. Events:
|
|
- ``index_updated``: partial index change (file create/modify/delete/move)
|
|
- ``index_reloaded``: full re-index completed
|
|
- ``vault_added``: new vault added dynamically
|
|
- ``vault_removed``: vault removed dynamically
|
|
"""
|
|
queue = await sse_manager.connect()
|
|
|
|
async def event_generator():
|
|
try:
|
|
# Send initial connection event
|
|
yield f"event: connected\ndata: {_json.dumps({'sse_clients': sse_manager.client_count})}\n\n"
|
|
while True:
|
|
try:
|
|
msg = await asyncio.wait_for(queue.get(), timeout=30.0)
|
|
yield f"event: {msg['event']}\ndata: {msg['data']}\n\n"
|
|
except asyncio.TimeoutError:
|
|
# Keepalive comment
|
|
yield ": keepalive\n\n"
|
|
except asyncio.CancelledError:
|
|
break
|
|
finally:
|
|
sse_manager.disconnect(queue)
|
|
|
|
return StreamingResponse(
|
|
event_generator(),
|
|
media_type="text/event-stream",
|
|
headers={
|
|
"Cache-Control": "no-cache",
|
|
"Connection": "keep-alive",
|
|
"X-Accel-Buffering": "no",
|
|
},
|
|
)
|
|
|
|
|
|
@router.websocket("/ws/collab/{vault_name}/{path:path}")
|
|
async def collab_websocket(websocket: WebSocket, vault_name: str, path: str):
|
|
"""Real-time collaborative editing over WebSocket (ROADMAP #62).
|
|
|
|
One *room* is created per ``vault::path``; all clients editing the same
|
|
file share Yjs/CRDT updates, awareness (cursors/selection) and a debounced
|
|
server-side persistence of the markdown content.
|
|
|
|
Authentication is performed manually (FastAPI ``Depends`` do not run for
|
|
WebSocket routes) and vault access is enforced per connection.
|
|
"""
|
|
from backend.services.errors import ServiceError
|
|
|
|
user = authenticate_websocket(websocket)
|
|
if user is None:
|
|
await websocket.close(code=4401)
|
|
return
|
|
|
|
if not check_vault_access(vault_name, user):
|
|
await websocket.close(code=4403)
|
|
return
|
|
|
|
try:
|
|
vault_root = get_vault_root(vault_name)
|
|
file_path = resolve_safe_path(vault_root, path)
|
|
except ServiceError:
|
|
await websocket.close(code=4404)
|
|
return
|
|
|
|
if not file_path.exists() or not file_path.is_file():
|
|
await websocket.close(code=4404)
|
|
return
|
|
|
|
await websocket.accept()
|
|
await collab_manager.connect(websocket, vault_name, path, file_path, user)
|