- A25 — 84 `except Exception: pass/…` → `logger.exception("<fonction>")`
(19 fichiers : api_v2 30, dashboard 10, board 7, sites 5, workspace 5,
api_v2_helpers 5, …) ; `logger` ajouté là où il manquait (api_v2_helpers,
sites + `import logging`)
- A25 critique — les `try` autour de `materialize_properties` supprimés dans
`create_collection_v2` ET `apply_db_template_v2` : un échec interrompt la
transaction au lieu de commiter une collection sans schéma
- test `test_collection_rollback_when_materialize_fails` (Bearer v2, monkeypatch
qui lève, assertions : RuntimeError + 0 ligne commitée)
- A21 partiel — `PRAGMA busy_timeout=5000` dans `get_conn()` (point d'entrée
unique) ; commentaire `ponytail:` : le wrapper async + les 510 call sites
restent à migrer module par module
- suite **1028/1028** · `ruff check app tests` OK
531 lines
21 KiB
Python
531 lines
21 KiB
Python
"""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()
|