Files
flowdeck/app/services/realtime_server.py
T
bruno 4038e9bdad
FlowDeck CI / lint (push) Successful in 52s
FlowDeck CI / test (push) Successful in 5m46s
FlowDeck CI / docker (push) Successful in 47s
fix(editor): repair block identity for template pages (duplicate/reordered lines)
Template-created pages (weekly report, project doc, meeting notes, to-do
list, ...) were persisted without block ids. The editor assigned ids
client-side, but the realtime room loaded the raw id-less content and sent
it back on 'sync'; applySync then merged id-less server blocks with local
blocks, producing data-bid='undefined' collisions and duplicated/shuffled
lines as soon as the user edited. Editing an empty page was unaffected
because the server state was empty.

Fixes:
- board.use_page_template: materialize unique block ids (recursively) when
  instantiating built-in or user templates.
- realtime_server: unique block_id() (uuid) + ensure_block_ids() on room
  load and on insert ops.
- editor: recursive ensureBlockIds() in init; gtTok() now restores
  [[fddate:...]] tokens (date chips survived as labels before).
- realtime client: applySync() normalizes ids, no longer appends unknown
  local blocks (duplication), and keeps local content when server is empty.

Tests: pytest (templates + realtime) and Playwright e2e covering to-do
list, weekly report date chip, Enter ordering and legacy id-less repair.
2026-09-14 11:45:44 -04:00

314 lines
10 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 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."""
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 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 = ensure_block_ids(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()