Files
flowdeck/app/services/realtime_server.py
T
bruno 3913f9f129
FlowDeck CI / test (push) Failing after 17s
FlowDeck CI / docker (push) Skipped
v5.13.0 Realtime : édition collaborative en direct
- Passerelle WebSocket WS /ws/pages/{page_id} (auth cookie, close 4401/4404), rooms par page en mémoire
- Protocole hello/sync/op/ack/title/sel/ping/peer_join/peer_leave ; version de page + stale => resync
- Merge des ops de blocs insert/update/delete/move, last-write-wins par bloc, ordre d'arrivée
- Client éditeur : diff local -> ops (debounce), application distante discret (LWW sur bloc focalisé), re-focus du bloc actif
- Présence (avatars topbar) + curseurs live (calque dédié, positions à l'édition/au scroll)
- Titre synchronisé (debounce) sans écraser le titre en cours d'édition
- Fallback polling 10 s si WS indisponible (adopté seulement sans brouillon local) + reconnexion auto
- Persistance debounce ~1 s (content/title) + flush à la déconnexion du dernier client
- CSP connect-src étendu à ws: ; peers sans données sensibles (id/login/full_name/couleur)
- 12 tests tests/test_realtime.py (auth, page absente, hello->sync, LWW 2 clients + persistance, présence, curseurs, titre, stale resync, apply_op/merge_ops) ; suite 298 verte (3 PDF pré-existants)
- Bump v5.3.0 (VERSION, main.py, CHANGELOG) ; ROADMAP v5.13.0 livré
2026-09-08 06:22:14 -04:00

294 lines
9.8 KiB
Python

"""FlowDeck — v5.13.0 Realtime: WebSocket gateway, présences, curseurs live,
merge des opérations de blocs (last-write-wins par bloc) + version de page.
Rooms in-memory (un seul worker uvicorn). Persistance en base (page.content)
avec debounce. Fallback polling côté client si le WS est indisponible.
"""
from __future__ import annotations
import asyncio
import json
import logging
from fastapi import WebSocket
from app.db import get_conn
logger = logging.getLogger(__name__)
COLORS = ["#2383E2", "#46A758", "#E5484D", "#F76B15", "#8E4EC6", "#12A594",
"#FFC53D", "#D6409F", "#0091FF", "#3E63DD", "#30A46C", "#FF3333"]
def color_for(uid: int) -> str:
return COLORS[(uid or 0) % len(COLORS)]
def block_id() -> str:
import time
return f"b{int(time.time()*1000)}"
def apply_op(blocks: list[dict], op: dict) -> list[dict]:
"""Apply one block op (insert/update/delete/move) — LWW par bloc."""
t = op.get("type")
if t == "insert":
blk = op.get("block") or {}
if not blk.get("id"):
blk = dict(blk)
blk["id"] = block_id()
idx = op.get("index")
if not isinstance(idx, int):
idx = len(blocks)
idx = max(0, min(idx, len(blocks)))
return blocks[:idx] + [blk] + blocks[idx:]
if t == "update":
nb = op.get("block") or {}
if not nb.get("id"):
return blocks
return [nb if b.get("id") == nb["id"] else b for b in blocks]
if t == "delete":
bid = op.get("id")
return [b for b in blocks if b.get("id") != bid]
if t == "move":
bid = op.get("id")
idx = op.get("index", 0) or 0
out = [b for b in blocks if b.get("id") != bid]
idx = max(0, min(idx, len(out)))
moved = next((b for b in blocks if b.get("id") == bid), None)
if moved is None:
return blocks
out.insert(idx, moved)
return out
return blocks
def merge_ops(blocks: list[dict], ops: list[dict]) -> list[dict]:
"""Apply a batch of ops sequentially (arrival order)."""
out = blocks
for op in ops or []:
out = apply_op(out, op)
return out
class Room:
__slots__ = ("page_id", "blocks", "title", "version", "conns",
"persist_task", "dirty")
def __init__(self, page_id: int):
self.page_id = page_id
self.blocks: list[dict] = []
self.title = ""
self.version = 0
self.conns: set["RTConn"] = set()
self.persist_task: asyncio.Task | None = None
self.dirty = False
class RTConn:
__slots__ = ("ws", "user", "page_id")
def __init__(self, ws: WebSocket, user: dict, page_id: int):
self.ws = ws
self.user = user
self.page_id = page_id
class RealtimeManager:
def __init__(self):
self._rooms: dict[int, Room] = {}
def room(self, page_id: int) -> Room:
return self._rooms.setdefault(page_id, Room(page_id))
@staticmethod
def _peer(user: dict) -> dict:
uid = user.get("id") or 0
return {
"id": uid,
"login": user.get("login", ""),
"full_name": user.get("full_name", "") or user.get("login", ""),
"color": color_for(uid),
}
async def load_room(self, room: Room) -> bool:
try:
with get_conn() as conn:
row = conn.execute(
"SELECT title, content, content_format FROM pages WHERE id=? AND deleted_at IS NULL",
(room.page_id,),
).fetchone()
except Exception as e:
logger.warning("realtime load failed: %s", e)
return False
if not row:
return False
room.title = row["title"] or ""
if (row["content_format"] or "") == "blocks" and row["content"]:
try:
room.blocks = json.loads(row["content"])
except Exception:
room.blocks = []
return True
async def connect(self, ws: WebSocket, page_id: int, user: dict) -> RTConn | None:
room = self.room(page_id)
if not room.conns and not await self.load_room(room):
await ws.close(code=4404)
return None
conn = RTConn(ws, user, page_id)
room.conns.add(conn)
me = self._peer(user)
peers = [self._peer(c.user) for c in room.conns if c is not conn]
await ws.send_json({"t": "welcome", "self": me,
"peers": peers, "color": me["color"]})
await ws.send_json({"t": "sync", "blocks": room.blocks,
"title": room.title, "version": room.version})
for c in room.conns:
if c is not conn:
try:
await c.ws.send_json({"t": "peer_join", "peer": me})
except Exception:
pass
return conn
async def disconnect(self, conn: RTConn):
room = self._rooms.get(conn.page_id)
if not room:
return
room.conns.discard(conn)
for c in room.conns:
try:
await c.ws.send_json({"t": "peer_leave",
"id": conn.user.get("id") or 0})
except Exception:
pass
if not room.conns:
await self.flush(room)
self._rooms.pop(conn.page_id, None)
async def flush(self, room: Room):
"""Write current room state to DB (sync, used on idle + disconnect)."""
if room.persist_task and not room.persist_task.done():
room.persist_task.cancel()
await self._persist(room)
async def _persist(self, room: Room):
try:
with get_conn() as conn:
conn.execute(
"UPDATE pages SET content=?, title=?, updated_at=CURRENT_TIMESTAMP WHERE id=?",
(json.dumps(room.blocks, ensure_ascii=False), room.title, room.page_id),
)
conn.commit()
room.dirty = False
except Exception as e:
logger.warning("realtime persist failed: %s", e)
def _schedule_persist(self, room: Room):
if room.persist_task and not room.persist_task.done():
return
room.dirty = True
async def _run():
try:
await asyncio.sleep(1.2)
await self._persist(room)
except asyncio.CancelledError:
pass
room.persist_task = asyncio.create_task(_run())
async def _broadcast(self, room: Room, msg: dict, exclude: RTConn | None = None):
for c in room.conns:
if c is exclude:
continue
try:
await c.ws.send_json(msg)
except Exception:
pass
async def handle(self, conn: RTConn, msg: dict):
room = self._rooms.get(conn.page_id)
if not room:
return
t = msg.get("t")
me = (conn.user.get("id") or 0)
if t == "hello":
try:
await conn.ws.send_json({"t": "sync", "blocks": room.blocks,
"title": room.title, "version": room.version})
except Exception:
pass
return
if t == "sync_req":
try:
await conn.ws.send_json({"t": "sync", "blocks": room.blocks,
"title": room.title, "version": room.version})
except Exception:
pass
return
if t == "ping":
try:
await conn.ws.send_json({"t": "pong"})
except Exception:
pass
return
if t == "op":
op = msg.get("op") or {}
client_v = msg.get("v", 0)
room.blocks = apply_op(room.blocks, op)
room.version += 1
self._schedule_persist(room)
stale = client_v < room.version - 1
await self._broadcast(room, {"t": "op", "op": op, "from": me,
"v": room.version}, exclude=conn)
try:
await conn.ws.send_json({"t": "ack", "v": room.version,
"stale": stale})
except Exception:
pass
if stale:
try:
await conn.ws.send_json({"t": "sync",
"blocks": room.blocks,
"title": room.title,
"version": room.version})
except Exception:
pass
return
if t == "title":
fmt = (msg.get("title") or "").strip()
if fmt and fmt != room.title:
room.title = fmt
room.version += 1
self._schedule_persist(room)
await self._broadcast(room, {"t": "title", "title": room.title,
"from": me, "v": room.version}, exclude=conn)
try:
await conn.ws.send_json({"t": "ack", "v": room.version,
"stale": False})
except Exception:
pass
return
if t == "sel":
await self._broadcast(room, {"t": "sel", "from": me,
"peer": self._peer(conn.user),
"block": msg.get("block"),
"offset": msg.get("offset", 0)}, exclude=conn)
return
async def room_state(self, page_id: int) -> dict:
room = self.room(page_id)
if not room.conns and not room.blocks:
await self.load_room(room)
return {"blocks": room.blocks, "title": room.title, "version": room.version}
manager = RealtimeManager()