"""FlowDeck — v6.4.0 Realtime: WebSocket gateway, présences, curseurs live, merge à trois versions des opérations de blocs (au-delà du last-write-wins) + version de page. Améliorations « production » par rapport à v5.13.0 : * **Conflits** — les mises à jour de bloc embarquent la ``base`` dont dérive la saisie du client ; le serveur fait un merge 3-voix champ-par-champ + texte (voir ``realtime_merge``) au lieu d'écraser le bloc entier. Les deux saisies disjointes survivent, les chevauchements réels retombent en LWW *par champ* avec drapeau de conflit renvoyé au client. * **Échelle** — chaque connexion possède une file sortante + une tâche writer dédiée ; le broadcast devient non bloquant (un client lent ne bloque plus la room), les mises à jour de curseur se coalescent (une seule par flush), et un client trop lent (file pleine) est déconnecté proprement (4413). * **Anti-flood** — budget d'opérations par connexion (fenêtre glissante). * **Fuites corrigées** — une room 4404 n'est plus enregistrée ; ``room_state`` ne crée plus d'objet None ; les rooms vides sont évacuées. * **Observabilité** — compteurs (rooms, conns, ops, merges, conflits, déconnexions lentes) exposés par ``stats()``. 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 import time from fastapi import WebSocket from app.db import get_conn from app.services.realtime_merge import merge_block_3way logger = logging.getLogger(__name__) COLORS = ["#2383E2", "#46A758", "#E5484D", "#F76B15", "#8E4EC6", "#12A594", "#FFC53D", "#D6409F", "#0091FF", "#3E63DD", "#30A46C", "#FF3333"] # File sortante maximale par connexion au-delà de laquelle le client est # considéré comme trop lent et déconnecté (évite qu'une room entière stagne). MAX_OUT_QUEUE = 512 # Nombre maximal d'opérations acceptées par connexion et par fenêtre (anti-flood). OP_WINDOW_SECONDS = 10.0 OP_WINDOW_MAX = 400 def color_for(uid: int) -> str: return COLORS[(uid or 0) % len(COLORS)] def block_id() -> str: import uuid return "b" + uuid.uuid4().hex[:12] def ensure_block_ids(blocks: list[dict]) -> list[dict]: """Assign unique ids to blocks missing one, recursively. Template-created pages may have been persisted without block ids; without them the room state is not addressable by ops and the editor ends up with ``data-bid="undefined"`` blocks (duplicated / reordered lines). """ if not isinstance(blocks, list): return blocks for b in blocks: if isinstance(b, dict): if not b.get("id"): b["id"] = block_id() if isinstance(b.get("children"), list): ensure_block_ids(b["children"]) return blocks def apply_op(blocks: list[dict], op: dict) -> list[dict]: """Apply one block op (insert/update/delete/move) — LWW par bloc. Conservé pour la compatibilité : tests et chemins sans ``base`` continuent de fonctionner. Le merge 3-voix vit dans ``RealtimeManager._apply``. """ t = op.get("type") if t == "insert": blk = op.get("block") or {} if not blk.get("id"): blk = dict(blk) blk["id"] = block_id() ensure_block_ids([blk]) 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 RTConn: """Une connexion WS : file sortante + tâche writer dédiée. Le broadcast ne fait que ``put_nowait`` dans la file ; c'est la tâche writer qui consomme et écrit sur le socket. Un client lent n'empêche donc jamais les autres membres de la room de recevoir les messages. """ __slots__ = ("ws", "user", "page_id", "out_q", "writer", "closed", "_ops_count", "_ops_window_start", "created_at") def __init__(self, ws: WebSocket, user: dict, page_id: int): self.ws = ws self.user = user self.page_id = page_id self.out_q: asyncio.Queue = asyncio.Queue(maxsize=MAX_OUT_QUEUE) self.writer: asyncio.Task | None = None self.closed = False self._ops_count = 0 self._ops_window_start = time.monotonic() self.created_at = time.monotonic() def op_budget_ok(self) -> bool: """Fenêtre glissante simple anti-flood d'opérations.""" now = time.monotonic() if now - self._ops_window_start > OP_WINDOW_SECONDS: self._ops_window_start = now self._ops_count = 0 self._ops_count += 1 return self._ops_count <= OP_WINDOW_MAX class Room: __slots__ = ("page_id", "blocks", "title", "version", "conns", "persist_task", "dirty", "merge_count", "conflict_count") 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 self.merge_count = 0 self.conflict_count = 0 class RealtimeManager: def __init__(self): self._rooms: dict[int, Room] = {} # compteurs globaux (observabilité) self.stat_ops = 0 self.stat_merges = 0 self.stat_conflicts = 0 self.stat_slow_disconnects = 0 self.stat_connections_total = 0 # ── rooms ──────────────────────────────────────────────────────────── 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: blocks = ensure_block_ids(json.loads(row["content"])) # v6.5.0: rooms serve server-resolved synced blocks so a # sync_req never returns a stale cached copy. from app.services.synced_blocks import resolve_synced_block room.blocks = resolve_synced_block(blocks) except Exception: room.blocks = [] return True # ── entrées / sorties ──────────────────────────────────────────────── async def connect(self, ws: WebSocket, page_id: int, user: dict) -> RTConn | None: # Ne JAMAIS enregistrer la room avant d'avoir validé l'existence de la # page : avant, un 4404 laissait une Room orpheline en mémoire pour # toujours (fuite). existing = self._rooms.get(page_id) room = existing if existing is not None else Room(page_id) if not room.conns and not await self.load_room(room): await ws.close(code=4404) return None # room non enregistrée → pas de fuite if existing is None: self._rooms[page_id] = room conn = RTConn(ws, user, page_id) self.stat_connections_total += 1 conn.writer = asyncio.create_task(self._writer(conn)) 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}) await self._broadcast(room, {"t": "peer_join", "peer": me}, exclude=conn) return conn async def disconnect(self, conn: RTConn): room = self._rooms.get(conn.page_id) if room is None: self._close_writer(conn) return room.conns.discard(conn) self._close_writer(conn) await self._broadcast(room, {"t": "peer_leave", "id": conn.user.get("id") or 0}) if not room.conns: await self.flush(room) self._rooms.pop(conn.page_id, None) def _close_writer(self, conn: RTConn): conn.closed = True t = conn.writer if t is not None and not t.done(): t.cancel() conn.writer = None async def _writer(self, conn: RTConn): """Tâche dédiée : consomme la file sortante et écrit sur le socket. Sert aussi de filet de coalescence : après chaque message consommé, les mises à jour de curseur (``sel``) déjà empilées sont réduites à la dernière (les curseurs n'ont pas besoin d'être ordonnés entre eux, seule la position la plus récente compte). """ try: while not conn.closed: msg = await conn.out_q.get() if msg is None: return await conn.ws.send_json(msg) # coalescence des curseurs en attente pending_sel = [] while not conn.out_q.empty(): nxt = conn.out_q.get_nowait() if nxt is None: return if isinstance(nxt, dict) and nxt.get("t") == "sel": pending_sel.append(nxt) else: await conn.ws.send_json(nxt) if pending_sel: await conn.ws.send_json(pending_sel[-1]) except asyncio.CancelledError: raise except Exception as e: # socket mort → on arrête proprement logger.debug("writer stopped: %s", e) conn.closed = True # ── persistance ────────────────────────────────────────────────────── 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()) # ── broadcast non bloquant ─────────────────────────────────────────── def _enqueue(self, conn: RTConn, msg: dict) -> bool: """Pose le message dans la file du client. Retourne False si trop lent.""" if conn.closed: return False try: conn.out_q.put_nowait(msg) return True except asyncio.QueueFull: return False async def _evict_slow(self, room: Room, conn: RTConn): """Client dépassé : on le déconnecte pour ne pas figer la room.""" self.stat_slow_disconnects += 1 logger.info("realtime: evicting slow client (user=%s page=%s)", conn.user.get("login"), conn.page_id) room.conns.discard(conn) self._close_writer(conn) try: await conn.ws.close(code=4413) except Exception: logger.exception("_evict_slow") async def _broadcast(self, room: Room, msg: dict, exclude: RTConn | None = None): """Enfile ``msg`` chez chaque membre — jamais d'attente sur le socket.""" slow: list[RTConn] = [] for c in list(room.conns): if c is exclude: continue if not self._enqueue(c, msg): slow.append(c) for c in slow: await self._evict_slow(room, c) async def _send(self, conn: RTConn, msg: dict): if not self._enqueue(conn, msg): room = self._rooms.get(conn.page_id) if room is not None: await self._evict_slow(room, conn) # ── protocole ──────────────────────────────────────────────────────── async def handle(self, conn: RTConn, msg: dict): room = self._rooms.get(conn.page_id) if room is None: return t = msg.get("t") me = (conn.user.get("id") or 0) if t == "hello": await self._send(conn, {"t": "sync", "blocks": room.blocks, "title": room.title, "version": room.version}) return if t == "sync_req": await self._send(conn, {"t": "sync", "blocks": room.blocks, "title": room.title, "version": room.version}) return if t == "ping": await self._send(conn, {"t": "pong"}) return if t == "op": if not conn.op_budget_ok(): # trop d'ops : on ignore silencieusement (le client resync) await self._send(conn, {"t": "ack", "v": room.version, "stale": True}) return op = msg.get("op") or {} client_v = msg.get("v", 0) result = self._apply(room, op) room.version += 1 self.stat_ops += 1 self._schedule_persist(room) stale = client_v < room.version - 1 # broadcast : on diffuse TOUJOURS le bloc final fusionné (pas la # proposition brute) pour que tous les clients convergent. out_op = dict(op) if result.get("merged") is not None: out_op["block"] = result["merged"] out_op["merged"] = True await self._broadcast(room, {"t": "op", "op": out_op, "from": me, "v": room.version, "conflict": result.get("conflict", False)}, exclude=conn) ack = {"t": "ack", "v": room.version, "stale": stale} if result.get("merged") is not None: ack["merged"] = result["merged"] ack["conflict"] = result.get("conflict", False) await self._send(conn, ack) if stale: await self._send(conn, {"t": "sync", "blocks": room.blocks, "title": room.title, "version": room.version}) 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) await self._send(conn, {"t": "ack", "v": room.version, "stale": False}) 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 def _apply(self, room: Room, op: dict) -> dict: """Applique une opération en mergeant à 3 voix si le client fournit une ``base``. Retourne ``{"merged": bloc|None, "conflict": bool}``. * ``update`` avec ``base`` → merge 3-voix (au-delà du LWW). * le reste (insert/delete/move, ou update sans base) → LWW historique. """ if op.get("type") != "update" or "base" not in op: room.blocks = apply_op(room.blocks, op) return {"merged": None, "conflict": False} incoming = op.get("block") or {} base = op.get("base") bid = incoming.get("id") if not bid: return {"merged": None, "conflict": False} current = next((b for b in room.blocks if b.get("id") == bid), None) if current is None: # bloc introuvable : LWW historique = pas d'update possible return {"merged": None, "conflict": False} merged, conflicts = merge_block_3way(base, current, incoming) room.merge_count += 1 conflict = bool(conflicts) if conflict: room.conflict_count += 1 self.stat_conflicts += 1 logger.debug("realtime conflict page=%s block=%s fields=%s", room.page_id, bid, conflicts) self.stat_merges += 1 room.blocks = [merged if b.get("id") == bid else b for b in room.blocks] return {"merged": merged, "conflict": conflict} # ── observabilité ──────────────────────────────────────────────────── async def room_state(self, page_id: int) -> dict: room = self._rooms.get(page_id) if room is None: # ne pas créer d'objet None : on charge dans une room jetable room = Room(page_id) if not await self.load_room(room): return {"blocks": [], "title": "", "version": 0} elif not room.conns and not room.blocks: await self.load_room(room) return {"blocks": room.blocks, "title": room.title, "version": room.version} def stats(self) -> dict: rooms = len(self._rooms) conns = sum(len(r.conns) for r in self._rooms.values()) return { "rooms": rooms, "connections": conns, "connections_total": self.stat_connections_total, "ops": self.stat_ops, "merges": self.stat_merges, "conflicts": self.stat_conflicts, "slow_disconnects": self.stat_slow_disconnects, "pages": [{"page_id": r.page_id, "conns": len(r.conns), "version": r.version, "merges": r.merge_count, "conflicts": r.conflict_count} for r in self._rooms.values()], } async def _broadcast_synced_to(self, pages: list[int], synced_id: int) -> None: """Broadcast a synced-block event to the given pages' open rooms.""" for pid in pages: room = self._rooms.get(pid) if room and room.conns: await self.load_room(room) await self._broadcast(room, {"t": "synced_update", "synced_id": synced_id, "version": room.version}) async def _propagate_synced(self, synced_id: int) -> None: """Broadcast a synced-block update to all rooms that reference it.""" from app.services.synced_blocks import page_ids_for_synced try: pages = page_ids_for_synced(synced_id) except Exception: return await self._broadcast_synced_to(pages, synced_id) manager = RealtimeManager()